61 using Handler = std::function<void(
const PeerId& from,
const std::string& topic, ByteView data)>;
71 int gossip_factor = 6;
72 std::chrono::milliseconds fanout_ttl{60000};
73 std::chrono::milliseconds heartbeat_interval{1000};
74 int history_length = 5;
75 int history_gossip = 3;
76 size_t seen_limit = 8192;
88 void publish(
const std::string& topic, ByteView data);
93 std::vector<PeerId>
mesh_peers(
const std::string& topic)
const;
106 std::unordered_set<PeerId, PeerId::Hash> subscribers;
107 std::unordered_set<PeerId, PeerId::Hash> mesh;
108 std::unordered_set<PeerId, PeerId::Hash> fanout;
109 std::chrono::steady_clock::time_point last_fanout{};
112 struct CachedMessage {
118 void on_new_peer(
const Peer& peer);
119 void on_peer_gone(
const PeerId&
id);
120 void on_gossip(
const Peer& peer, ByteView payload);
122 void recv_subscription(
const PeerId& from,
const std::string& topic,
bool subscribe);
123 void recv_graft(
const PeerId& from,
const std::string& topic);
124 void recv_prune(
const PeerId& from,
const std::string& topic);
125 void recv_publish(
const PeerId& from, ByteView frame,
const PeerId& origin, uint64_t seqno,
126 const std::string& topic, ByteView data);
127 void recv_ihave(
const PeerId& from,
const std::string& topic,
const std::vector<std::string>& ids);
128 void recv_iwant(
const PeerId& from,
const std::vector<std::string>& ids);
131 void heartbeat_loop();
137 using CtrlList = std::vector<std::pair<PeerId, std::string>>;
138 void maintain_mesh_locked(
const std::string& topic, CtrlList& grafts, CtrlList& prunes);
141 void send_ctrl(
const PeerId& to, uint8_t op,
const std::string& topic);
142 void broadcast_ctrl(uint8_t op,
const std::string& topic);
143 void deliver_local(
const PeerId& from,
const std::string& topic, ByteView data);
144 Handler handler_for(
const std::string& topic)
const;
145 ValidationResult validate(
const PeerId& from,
const std::string& topic, ByteView data);
146 bool mark_seen(
const std::string&
id);
147 void cache_message(
const std::string&
id,
const std::string& topic,
const Bytes& frame);
148 std::vector<PeerId> random_sample(std::vector<PeerId> in,
int n);
150 PeerNetwork* network_ =
nullptr;
153 mutable std::mutex mutex_;
155 std::unordered_map<std::string, Handler> subscriptions_;
156 std::unordered_map<std::string, Topic> topics_;
157 std::unordered_map<std::string, Validator> validators_;
158 Validator global_validator_;
160 mutable std::mutex mcache_mutex_;
161 std::unordered_map<std::string, CachedMessage> mcache_;
162 std::deque<std::vector<std::string>> history_;
163 std::unordered_set<std::string> seen_;
164 std::deque<std::string> seen_order_;
166 std::mutex rng_mutex_;
169 std::atomic<bool> running_{
false};
170 std::thread heartbeat_thread_;
171 std::mutex hb_mutex_;
172 std::condition_variable hb_cv_;
void publish(const std::string &topic, ByteView data)
Publish data on topic: along the mesh if we are subscribed, else via fanout.