Back to Site
Loading...
Searching...
No Matches
pubsub.h
Go to the documentation of this file.
1#pragma once
2
31#include "util/rats_export.h"
32#include "node/peer_network.h"
33#include "peer/peer.h"
34#include "core/bytes.h"
35#include "peer/peer_id.h"
36
37#include <atomic>
38#include <chrono>
39#include <condition_variable>
40#include <cstdint>
41#include <deque>
42#include <functional>
43#include <mutex>
44#include <random>
45#include <string>
46#include <thread>
47#include <unordered_map>
48#include <unordered_set>
49#include <utility>
50#include <vector>
51
52namespace librats {
53
58
59class RATS_API PubSub final : public Subsystem {
60public:
61 using Handler = std::function<void(const PeerId& from, const std::string& topic, ByteView data)>;
62 using Validator = std::function<ValidationResult(const PeerId& from, const std::string& topic, ByteView data)>;
63
66 struct Config {
67 int mesh_target = 6;
68 int mesh_low = 4;
69 int mesh_high = 12;
70 int fanout_size = 6;
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;
77 };
78
80 explicit PubSub(Config config);
81 ~PubSub() override;
82
84 void subscribe(const std::string& topic, Handler handler);
85 void unsubscribe(const std::string& topic);
86
88 void publish(const std::string& topic, ByteView data);
89
90 bool is_subscribed(const std::string& topic) const;
91 std::vector<std::string> subscribed_topics() const;
92 std::vector<PeerId> peers_for_topic(const std::string& topic) const;
93 std::vector<PeerId> mesh_peers(const std::string& topic) const;
94
97 void set_validator(const std::string& topic, Validator validator);
98
99 // Subsystem.
100 void attach(NodeContext& ctx) override;
101 void start() override;
102 void stop() override;
103
104private:
105 struct Topic {
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{};
110 };
111
112 struct CachedMessage {
113 std::string topic;
114 Bytes frame;
115 };
116
117 // Inbound dispatch (all run on a reactor thread).
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);
121
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);
129
130 // Heartbeat (background thread).
131 void heartbeat_loop();
132 void do_heartbeat();
133
137 using CtrlList = std::vector<std::pair<PeerId, std::string>>;
138 void maintain_mesh_locked(const std::string& topic, CtrlList& grafts, CtrlList& prunes);
139
140 // Helpers.
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);
149
150 PeerNetwork* network_ = nullptr;
151 Config config_;
152
153 mutable std::mutex mutex_;
154 uint64_t seqno_ = 0;
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_;
159
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_;
165
166 std::mutex rng_mutex_;
167 std::mt19937 rng_;
168
169 std::atomic<bool> running_{false};
170 std::thread heartbeat_thread_;
171 std::mutex hb_mutex_;
172 std::condition_variable hb_cv_;
173};
174
175} // namespace librats
void subscribe(const std::string &topic, Handler handler)
Subscribe to a topic and deliver matching messages to handler.
void publish(const std::string &topic, ByteView data)
Publish data on topic: along the mesh if we are subscribed, else via fanout.
std::function< void(const PeerId &from, const std::string &topic, ByteView data)> Handler
Definition pubsub.h:61
PubSub(Config config)
void attach(NodeContext &ctx) override
void unsubscribe(const std::string &topic)
~PubSub() override
std::function< ValidationResult(const PeerId &from, const std::string &topic, ByteView data)> Validator
Definition pubsub.h:62
std::vector< PeerId > peers_for_topic(const std::string &topic) const
known subscribers
void stop() override
stop and join it
void start() override
launch the heartbeat thread
std::vector< PeerId > mesh_peers(const std::string &topic) const
our mesh for topic
bool is_subscribed(const std::string &topic) const
void set_validator(const std::string &topic, Validator validator)
Gate inbound messages for a topic; topic == "" installs a global validator used when no per-topic val...
std::vector< std::string > subscribed_topics() const
A pluggable network subsystem.
Definition node.h:66
ValidationResult
Outcome of validating an inbound published message before it is delivered or forwarded.
Definition pubsub.h:57
A lightweight handle to a connected peer.
Self-certifying peer identity.
The narrow contract a subsystem needs from the node — and nothing more.
GossipSub tuning.
Definition pubsub.h:66