Use message timestamp for filtering
Signed-off-by: Alexis Sellier <alexis@radicle.xyz>
This commit is contained in:
parent
e35ec2f715
commit
af06ad6451
|
|
@ -453,7 +453,12 @@ where
|
||||||
let remote = repo.remote(&node).unwrap();
|
let remote = repo.remote(&node).unwrap();
|
||||||
let peers = self.sessions.negotiated().map(|(_, p)| p);
|
let peers = self.sessions.negotiated().map(|(_, p)| p);
|
||||||
let refs = remote.refs.into();
|
let refs = remote.refs.into();
|
||||||
let msg = AnnouncementMessage::from(RefsAnnouncement { id, refs });
|
let timestamp = self.clock.timestamp();
|
||||||
|
let msg = AnnouncementMessage::from(RefsAnnouncement {
|
||||||
|
id,
|
||||||
|
refs,
|
||||||
|
timestamp,
|
||||||
|
});
|
||||||
let ann = msg.signed(&self.signer);
|
let ann = msg.signed(&self.signer);
|
||||||
|
|
||||||
self.reactor.broadcast(ann, peers);
|
self.reactor.broadcast(ann, peers);
|
||||||
|
|
@ -708,7 +713,7 @@ where
|
||||||
|
|
||||||
// Returning true here means that the message should be relayed.
|
// Returning true here means that the message should be relayed.
|
||||||
if self.handle_announcement(&git, &ann)? {
|
if self.handle_announcement(&git, &ann)? {
|
||||||
self.gossip.received(ann.clone(), self.clock.timestamp());
|
self.gossip.received(ann.clone(), ann.message.timestamp());
|
||||||
return Ok(Some(ann));
|
return Ok(Some(ann));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -173,14 +173,17 @@ pub struct RefsAnnouncement {
|
||||||
pub id: Id,
|
pub id: Id,
|
||||||
/// Updated refs.
|
/// Updated refs.
|
||||||
pub refs: Refs,
|
pub refs: Refs,
|
||||||
// TODO: Add timestamp
|
/// Time of announcement.
|
||||||
|
pub timestamp: Timestamp,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Node announcing its inventory to the network.
|
/// Node announcing its inventory to the network.
|
||||||
/// This should be the whole inventory every time.
|
/// This should be the whole inventory every time.
|
||||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
pub struct InventoryAnnouncement {
|
pub struct InventoryAnnouncement {
|
||||||
|
/// Node inventory.
|
||||||
pub inventory: Vec<Id>,
|
pub inventory: Vec<Id>,
|
||||||
|
/// Time of announcement.
|
||||||
pub timestamp: Timestamp,
|
pub timestamp: Timestamp,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -207,6 +210,14 @@ impl AnnouncementMessage {
|
||||||
signature,
|
signature,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn timestamp(&self) -> Timestamp {
|
||||||
|
match self {
|
||||||
|
Self::Inventory(InventoryAnnouncement { timestamp, .. }) => *timestamp,
|
||||||
|
Self::Refs(RefsAnnouncement { timestamp, .. }) => *timestamp,
|
||||||
|
Self::Node(NodeAnnouncement { timestamp, .. }) => *timestamp,
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<NodeAnnouncement> for AnnouncementMessage {
|
impl From<NodeAnnouncement> for AnnouncementMessage {
|
||||||
|
|
@ -368,7 +379,12 @@ mod tests {
|
||||||
#[quickcheck]
|
#[quickcheck]
|
||||||
fn prop_refs_announcement_signing(id: Id, refs: Refs) {
|
fn prop_refs_announcement_signing(id: Id, refs: Refs) {
|
||||||
let signer = MockSigner::new(&mut fastrand::Rng::new());
|
let signer = MockSigner::new(&mut fastrand::Rng::new());
|
||||||
let message = AnnouncementMessage::Refs(RefsAnnouncement { id, refs });
|
let timestamp = 0;
|
||||||
|
let message = AnnouncementMessage::Refs(RefsAnnouncement {
|
||||||
|
id,
|
||||||
|
refs,
|
||||||
|
timestamp,
|
||||||
|
});
|
||||||
let ann = message.signed(&signer);
|
let ann = message.signed(&signer);
|
||||||
|
|
||||||
assert!(ann.verify());
|
assert!(ann.verify());
|
||||||
|
|
|
||||||
|
|
@ -61,6 +61,7 @@ impl Arbitrary for Message {
|
||||||
message: RefsAnnouncement {
|
message: RefsAnnouncement {
|
||||||
id: Id::arbitrary(g),
|
id: Id::arbitrary(g),
|
||||||
refs: Refs::arbitrary(g),
|
refs: Refs::arbitrary(g),
|
||||||
|
timestamp: Timestamp::arbitrary(g),
|
||||||
}
|
}
|
||||||
.into(),
|
.into(),
|
||||||
signature: crypto::Signature::from(ByteArray::<64>::arbitrary(g).into_inner()),
|
signature: crypto::Signature::from(ByteArray::<64>::arbitrary(g).into_inner()),
|
||||||
|
|
|
||||||
|
|
@ -12,8 +12,17 @@ pub fn messages(count: usize, now: LocalTime, delta: LocalDuration) -> Vec<Messa
|
||||||
|
|
||||||
for _ in 0..count {
|
for _ in 0..count {
|
||||||
let signer = MockSigner::new(&mut rng);
|
let signer = MockSigner::new(&mut rng);
|
||||||
let delta = LocalDuration::from_secs(rng.u64(0..delta.as_secs()));
|
let time = if delta == LocalDuration::from_secs(0) {
|
||||||
let time = if rng.bool() { now + delta } else { now - delta };
|
now
|
||||||
|
} else {
|
||||||
|
let delta = LocalDuration::from_secs(rng.u64(0..delta.as_secs()));
|
||||||
|
|
||||||
|
if rng.bool() {
|
||||||
|
now + delta
|
||||||
|
} else {
|
||||||
|
now - delta
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
msgs.push(Message::inventory(
|
msgs.push(Message::inventory(
|
||||||
InventoryAnnouncement {
|
InventoryAnnouncement {
|
||||||
|
|
|
||||||
|
|
@ -5,7 +5,7 @@ use crossbeam_channel as chan;
|
||||||
use nakamoto_net as nakamoto;
|
use nakamoto_net as nakamoto;
|
||||||
|
|
||||||
use crate::collections::{HashMap, HashSet};
|
use crate::collections::{HashMap, HashSet};
|
||||||
use crate::prelude::Timestamp;
|
use crate::prelude::{LocalDuration, Timestamp};
|
||||||
use crate::service::config::*;
|
use crate::service::config::*;
|
||||||
use crate::service::filter::Filter;
|
use crate::service::filter::Filter;
|
||||||
use crate::service::message::*;
|
use crate::service::message::*;
|
||||||
|
|
@ -223,7 +223,6 @@ fn test_gossip_rebroadcast() {
|
||||||
|
|
||||||
let received = test::gossip::messages(6, alice.local_time(), MAX_TIME_DELTA);
|
let received = test::gossip::messages(6, alice.local_time(), MAX_TIME_DELTA);
|
||||||
for msg in received.iter().cloned() {
|
for msg in received.iter().cloned() {
|
||||||
// TODO: Test with elapsed time
|
|
||||||
alice.receive(&bob.addr(), msg);
|
alice.receive(&bob.addr(), msg);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -241,6 +240,45 @@ fn test_gossip_rebroadcast() {
|
||||||
assert_eq!(relayed, received);
|
assert_eq!(relayed, received);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_gossip_rebroadcast_timestamp_filtered() {
|
||||||
|
let mut alice = Peer::new("alice", [7, 7, 7, 7], MockStorage::empty());
|
||||||
|
let bob = Peer::new("bob", [8, 8, 8, 8], MockStorage::empty());
|
||||||
|
let eve = Peer::new("eve", [9, 9, 9, 9], MockStorage::empty());
|
||||||
|
|
||||||
|
alice.connect_to(&bob);
|
||||||
|
|
||||||
|
let delta = LocalDuration::from_mins(10);
|
||||||
|
let first = test::gossip::messages(3, alice.local_time() - delta, LocalDuration::from_secs(0));
|
||||||
|
let second = test::gossip::messages(3, alice.local_time(), LocalDuration::from_secs(0));
|
||||||
|
let third = test::gossip::messages(3, alice.local_time() + delta, LocalDuration::from_secs(0));
|
||||||
|
|
||||||
|
// Alice receives three batches of messages.
|
||||||
|
for msg in first
|
||||||
|
.iter()
|
||||||
|
.chain(second.iter())
|
||||||
|
.chain(third.iter())
|
||||||
|
.cloned()
|
||||||
|
{
|
||||||
|
alice.receive(&bob.addr(), msg);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Eve subscribes to messages within the period of the second batch only.
|
||||||
|
alice.connect_from(&eve);
|
||||||
|
alice.receive(
|
||||||
|
&eve.addr(),
|
||||||
|
Message::Subscribe(Subscribe {
|
||||||
|
filter: Filter::default(),
|
||||||
|
since: alice.local_time().as_secs(),
|
||||||
|
until: (alice.local_time() + delta).as_secs(),
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
|
||||||
|
let relayed = alice.messages(&eve.addr()).collect::<Vec<_>>();
|
||||||
|
assert_eq!(relayed.len(), second.len());
|
||||||
|
assert_eq!(relayed, second);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_inventory_relay() {
|
fn test_inventory_relay() {
|
||||||
// Topology is eve <-> alice <-> bob
|
// Topology is eve <-> alice <-> bob
|
||||||
|
|
|
||||||
|
|
@ -100,6 +100,7 @@ impl wire::Encode for RefsAnnouncement {
|
||||||
|
|
||||||
n += self.id.encode(writer)?;
|
n += self.id.encode(writer)?;
|
||||||
n += self.refs.encode(writer)?;
|
n += self.refs.encode(writer)?;
|
||||||
|
n += self.timestamp.encode(writer)?;
|
||||||
|
|
||||||
Ok(n)
|
Ok(n)
|
||||||
}
|
}
|
||||||
|
|
@ -109,8 +110,13 @@ impl wire::Decode for RefsAnnouncement {
|
||||||
fn decode<R: std::io::Read + ?Sized>(reader: &mut R) -> Result<Self, wire::Error> {
|
fn decode<R: std::io::Read + ?Sized>(reader: &mut R) -> Result<Self, wire::Error> {
|
||||||
let id = Id::decode(reader)?;
|
let id = Id::decode(reader)?;
|
||||||
let refs = Refs::decode(reader)?;
|
let refs = Refs::decode(reader)?;
|
||||||
|
let timestamp = Timestamp::decode(reader)?;
|
||||||
|
|
||||||
Ok(Self { id, refs })
|
Ok(Self {
|
||||||
|
id,
|
||||||
|
refs,
|
||||||
|
timestamp,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue