use crate::rpc_proto;
use crate::TopicHash;
use tet_libp2p_core::PeerId;
use std::fmt;
use std::fmt::Debug;
#[derive(Debug)]
pub enum MessageAcceptance {
Accept,
Reject,
Ignore,
}
macro_rules! declare_message_id_type {
($name: ident, $name_string: expr) => {
#[derive(Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct $name(pub Vec<u8>);
impl $name {
pub fn new(value: &[u8]) -> Self {
Self(value.to_vec())
}
}
impl<T: Into<Vec<u8>>> From<T> for $name {
fn from(value: T) -> Self {
Self(value.into())
}
}
impl std::fmt::Display for $name {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", hex_fmt::HexFmt(&self.0))
}
}
impl std::fmt::Debug for $name {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}({})", $name_string, hex_fmt::HexFmt(&self.0))
}
}
};
}
declare_message_id_type!(MessageId, "MessageId");
declare_message_id_type!(FastMessageId, "FastMessageId");
#[derive(Debug, Clone, PartialEq)]
pub enum PeerKind {
Gossipsubv1_1,
Gossipsub,
Floodsub,
NotSupported,
}
#[derive(Clone, PartialEq, Eq, Hash, Debug)]
pub struct RawGossipsubMessage {
pub source: Option<PeerId>,
pub data: Vec<u8>,
pub sequence_number: Option<u64>,
pub topic: TopicHash,
pub signature: Option<Vec<u8>>,
pub key: Option<Vec<u8>>,
pub validated: bool,
}
#[derive(Clone, PartialEq, Eq, Hash)]
pub struct GossipsubMessage {
pub source: Option<PeerId>,
pub data: Vec<u8>,
pub sequence_number: Option<u64>,
pub topic: TopicHash,
}
impl fmt::Debug for GossipsubMessage {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("GossipsubMessage")
.field(
"data",
&format_args!("{:<20}", &hex_fmt::HexFmt(&self.data)),
)
.field("source", &self.source)
.field("sequence_number", &self.sequence_number)
.field("topic", &self.topic)
.finish()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct GossipsubSubscription {
pub action: GossipsubSubscriptionAction,
pub topic_hash: TopicHash,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub enum GossipsubSubscriptionAction {
Subscribe,
Unsubscribe,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct PeerInfo {
pub peer_id: Option<PeerId>,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub enum GossipsubControlAction {
IHave {
topic_hash: TopicHash,
message_ids: Vec<MessageId>,
},
IWant {
message_ids: Vec<MessageId>,
},
Graft {
topic_hash: TopicHash,
},
Prune {
topic_hash: TopicHash,
peers: Vec<PeerInfo>,
backoff: Option<u64>,
},
}
#[derive(Clone, PartialEq, Eq, Hash)]
pub struct GossipsubRpc {
pub messages: Vec<RawGossipsubMessage>,
pub subscriptions: Vec<GossipsubSubscription>,
pub control_msgs: Vec<GossipsubControlAction>,
}
impl GossipsubRpc {
pub fn into_protobuf(self) -> rpc_proto::Rpc {
self.into()
}
}
impl Into<rpc_proto::Rpc> for GossipsubRpc {
fn into(self) -> rpc_proto::Rpc {
let mut publish = Vec::new();
for message in self.messages.into_iter() {
let message = rpc_proto::Message {
from: message.source.map(|m| m.to_bytes()),
data: Some(message.data),
seqno: message.sequence_number.map(|s| s.to_be_bytes().to_vec()),
topic: TopicHash::into_string(message.topic),
signature: message.signature,
key: message.key,
};
publish.push(message);
}
let subscriptions = self
.subscriptions
.into_iter()
.map(|sub| rpc_proto::rpc::SubOpts {
subscribe: Some(sub.action == GossipsubSubscriptionAction::Subscribe),
topic_id: Some(sub.topic_hash.into_string()),
})
.collect::<Vec<_>>();
let mut control = rpc_proto::ControlMessage {
ihave: Vec::new(),
iwant: Vec::new(),
graft: Vec::new(),
prune: Vec::new(),
};
let empty_control_msg = self.control_msgs.is_empty();
for action in self.control_msgs {
match action {
GossipsubControlAction::IHave {
topic_hash,
message_ids,
} => {
let rpc_ihave = rpc_proto::ControlIHave {
topic_id: Some(topic_hash.into_string()),
message_ids: message_ids.into_iter().map(|msg_id| msg_id.0).collect(),
};
control.ihave.push(rpc_ihave);
}
GossipsubControlAction::IWant { message_ids } => {
let rpc_iwant = rpc_proto::ControlIWant {
message_ids: message_ids.into_iter().map(|msg_id| msg_id.0).collect(),
};
control.iwant.push(rpc_iwant);
}
GossipsubControlAction::Graft { topic_hash } => {
let rpc_graft = rpc_proto::ControlGraft {
topic_id: Some(topic_hash.into_string()),
};
control.graft.push(rpc_graft);
}
GossipsubControlAction::Prune {
topic_hash,
peers,
backoff,
} => {
let rpc_prune = rpc_proto::ControlPrune {
topic_id: Some(topic_hash.into_string()),
peers: peers
.into_iter()
.map(|info| rpc_proto::PeerInfo {
peer_id: info.peer_id.map(|id| id.to_bytes()),
signed_peer_record: None,
})
.collect(),
backoff,
};
control.prune.push(rpc_prune);
}
}
}
rpc_proto::Rpc {
subscriptions,
publish,
control: if empty_control_msg {
None
} else {
Some(control)
},
}
}
}
impl fmt::Debug for GossipsubRpc {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let mut b = f.debug_struct("GossipsubRpc");
if !self.messages.is_empty() {
b.field("messages", &self.messages);
}
if !self.subscriptions.is_empty() {
b.field("subscriptions", &self.subscriptions);
}
if !self.control_msgs.is_empty() {
b.field("control_msgs", &self.control_msgs);
}
b.finish()
}
}
impl PeerKind {
pub fn as_static_ref(&self) -> &'static str {
match self {
Self::NotSupported => "Not Supported",
Self::Floodsub => "Floodsub",
Self::Gossipsub => "Gossipsub v1.0",
Self::Gossipsubv1_1 => "Gossipsub v1.1",
}
}
}
impl AsRef<str> for PeerKind {
fn as_ref(&self) -> &str {
self.as_static_ref()
}
}
impl fmt::Display for PeerKind {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_ref())
}
}