node: Handle announcement commands

Signed-off-by: Alexis Sellier <self@cloudhead.io>
This commit is contained in:
Alexis Sellier 2022-09-01 18:47:34 +02:00
parent 400ba8d935
commit ff2bd185e7
No known key found for this signature in database
3 changed files with 43 additions and 21 deletions

View File

@ -208,7 +208,7 @@ impl<T: ReadStorage + WriteStorage, S: address_book::Store, G: crypto::Signer> P
/// Announce our inventory to all connected peers. /// Announce our inventory to all connected peers.
fn announce_inventory(&mut self) -> Result<(), storage::Error> { fn announce_inventory(&mut self) -> Result<(), storage::Error> {
let inv = Message::inventory(&mut self.context)?; let inv = Message::inventory(&self.context)?;
for addr in self.peers.negotiated().map(|(_, p)| p.addr) { for addr in self.peers.negotiated().map(|(_, p)| p.addr) {
self.context.write(addr, inv.clone()); self.context.write(addr, inv.clone());
@ -335,8 +335,11 @@ where
}) })
.unwrap(); .unwrap();
} }
Command::AnnounceInventory(_proj) => { Command::AnnounceInventory(proj) => {
todo!() let peers = self.peers.negotiated().map(|(_, p)| p.addr);
self.context
.broadcast(Message::InventoryUpdate { inv: vec![proj] }, peers);
} }
} }
} }
@ -470,10 +473,9 @@ where
let peers = negotiated let peers = negotiated
.iter() .iter()
.filter(|(ip, _)| *ip != peer.ip()) .filter(|(ip, _)| *ip != peer.ip())
.map(|(_, addr)| *addr) .map(|(_, addr)| *addr);
.collect::<Vec<_>>();
self.context.broadcast(msg, &peers); self.context.broadcast(msg, peers);
} }
Err(err) => { Err(err) => {
self.context self.context
@ -603,7 +605,23 @@ where
.entry(proj_id.clone()) .entry(proj_id.clone())
.or_insert_with(|| HashSet::with_hasher(self.rng.clone().into())); .or_insert_with(|| HashSet::with_hasher(self.rng.clone().into()));
// TODO: Fire an event on routing update.
if inventory.insert(from) && self.config.is_tracking(proj_id) {
self.fetch(proj_id, remote);
}
}
}
/// Process a peer inventory update announcement by (maybe) fetching.
fn process_inventory_update(&mut self, inventory: &Inventory, _from: NodeId, remote: &Url) {
for proj_id in inventory {
if self.config.is_tracking(proj_id) { if self.config.is_tracking(proj_id) {
self.fetch(proj_id, remote);
}
}
}
fn fetch(&mut self, proj_id: &ProjId, remote: &Url) {
// TODO: Verify refs before adding them to storage. // TODO: Verify refs before adding them to storage.
let mut repo = self.storage.repository(proj_id).unwrap(); let mut repo = self.storage.repository(proj_id).unwrap();
repo.fetch(&Url { repo.fetch(&Url {
@ -613,11 +631,6 @@ where
.unwrap(); .unwrap();
} }
// TODO: Fire an event on routing update.
inventory.insert(from);
}
}
/// Disconnect a peer. /// Disconnect a peer.
fn disconnect(&mut self, addr: net::SocketAddr, reason: DisconnectReason) { fn disconnect(&mut self, addr: net::SocketAddr, reason: DisconnectReason) {
self.io.push_back(Io::Disconnect(addr, reason)); self.io.push_back(Io::Disconnect(addr, reason));
@ -658,9 +671,9 @@ impl<S, T, G> Context<S, T, G> {
} }
/// Broadcast a message to a list of peers. /// Broadcast a message to a list of peers.
fn broadcast(&mut self, msg: Message, peers: &[net::SocketAddr]) { fn broadcast(&mut self, msg: Message, peers: impl IntoIterator<Item = net::SocketAddr>) {
for peer in peers { for peer in peers {
self.write(*peer, msg.clone()); self.write(peer, msg.clone());
} }
} }
} }

View File

@ -86,7 +86,9 @@ pub enum Message {
announcement: NodeAnnouncement, announcement: NodeAnnouncement,
}, },
/// Get a peer's inventory. /// Get a peer's inventory.
GetInventory { ids: Vec<ProjId> }, GetInventory {
ids: Vec<ProjId>,
},
/// Send our inventory to a peer. Sent in response to [`Message::GetInventory`]. /// Send our inventory to a peer. Sent in response to [`Message::GetInventory`].
/// Nb. This should be the whole inventory, not a partial update. /// Nb. This should be the whole inventory, not a partial update.
Inventory { Inventory {
@ -96,6 +98,9 @@ pub enum Message {
/// are the originator, only when relaying. /// are the originator, only when relaying.
origin: Option<NodeId>, origin: Option<NodeId>,
}, },
InventoryUpdate {
inv: Vec<ProjId>,
},
} }
impl Message { impl Message {
@ -119,7 +124,7 @@ impl Message {
} }
} }
pub fn inventory<S, T, G>(ctx: &mut Context<S, T, G>) -> Result<Self, storage::Error> pub fn inventory<S, T, G>(ctx: &Context<S, T, G>) -> Result<Self, storage::Error>
where where
T: storage::ReadStorage, T: storage::ReadStorage,
{ {

View File

@ -194,6 +194,10 @@ impl Peer {
})); }));
} }
} }
(PeerState::Negotiated { id, git, .. }, Message::InventoryUpdate { inv }) => {
// TODO: Buffer/throttle fetches.
ctx.process_inventory_update(&inv, *id, git);
}
( (
PeerState::Negotiated { .. }, PeerState::Negotiated { .. },
Message::Node { Message::Node {