sgns::crdt::PubSubBroadcasterExt¶
Extended PubSub broadcaster that integrates with a CRDT datastore and Graphsync DAG syncer. More...
#include <pubsub_broadcaster_ext.hpp>
Inherits from sgns::crdt::Broadcaster, std::enable_shared_from_this< PubSubBroadcasterExt >
Public Types¶
| Name | |
|---|---|
| using sgns::ipfs_pubsub::GossipPubSub | GossipPubSub |
Public Functions¶
| Name | |
|---|---|
| ~PubSubBroadcasterExt() | |
| virtual outcome::result< void > | Broadcast(const base::Buffer & buff, std::string topic, boost::optional< libp2p::peer::PeerInfo > peerInfo =boost::none) override Sends the given buffer as a broadcast to peers. |
| virtual outcome::result< base::Buffer > | Next() override Retrieves the next incoming broadcast payload. |
| virtual void | WaitForNext(std::chrono::milliseconds timeout) override Blocks until a message is queued or timeout elapses. |
| virtual void | CancelWait() override Releases any waiter and makes later waits return at once. |
| void | Start() Subscribes to all configured topics and starts message processing. Must be called before using Next() to receive incoming messages. |
| outcome::result< void > | AddBroadcastTopic(const std::string & topicName) Adds a new topic by name. |
| void | AddListenTopic(std::string topic) Subscribe to a given topic and store its future. |
| virtual bool | HasTopic(const std::string & topic) override Checks whether the given topic is already registered. |
| virtual std::shared_ptr< void > | GetDagSyncer() const override Get the underlying GraphsyncDAGSyncer instance. |
| void | Stop() |
| void | SetMembershipFilter(std::function< bool(const libp2p::peer::PeerId &)> filter) Installs (or replaces) the private-network membership filter consulted by OnMessage for EVERY inbound gossip message. |
| bool | HasMembershipFilter() const Reports whether a membership filter is currently installed. |
| void | ClearMembershipFilter() Removes the membership filter, restoring public pass-through ingest (teardown counterpart of SetMembershipFilter). |
| void | SetGossipSigningKey(std::shared_ptr< const libp2p::crypto::KeyPair > key) Installs the gossip host keypair used to SEAL private-network publishes (CR-G01 publisher side). |
| bool | HasGossipSigningKey() const Reports whether a gossip signing key is currently installed. |
| bool | AddSingleCIDInfo(const std::string & cid, const std::string peer_id, const std::string address) |
| std::shared_ptr< PubSubBroadcasterExt > | New(std::shared_ptr< sgns::crdt::GraphsyncDAGSyncer > dagSyncer, std::shared_ptr< GossipPubSub > pubSub) Factory method to create a broadcaster for multiple topics. |
Additional inherited members¶
Public Types inherited from sgns::crdt::Broadcaster
| Name | |
|---|---|
| enum class | ErrorCode |
Public Functions inherited from sgns::crdt::Broadcaster
| Name | |
|---|---|
| virtual | ~Broadcaster() =default |
Detailed Description¶
Extended PubSub broadcaster that integrates with a CRDT datastore and Graphsync DAG syncer.
Manages multiple gossip topics, broadcasting messages and processing incoming payloads.
Public Types Documentation¶
using GossipPubSub¶
Public Functions Documentation¶
function ~PubSubBroadcasterExt¶
function Broadcast¶
virtual outcome::result< void > Broadcast(
const base::Buffer & buff,
std::string topic,
boost::optional< libp2p::peer::PeerInfo > peerInfo =boost::none
) override
Sends the given buffer as a broadcast to peers.
Parameters:
- buff Buffer containing the data to broadcast.
- topic Topic to broadcast to.
- peerInfo Optional peer info to avoid repeated GetPeerInfo calls.
Return: outcome::success on successful publish, or outcome::failure on error.
Reimplements: sgns::crdt::Broadcaster::Broadcast
function Next¶
Retrieves the next incoming broadcast payload.
Return: buffer value or outcome::failure on error
Reimplements: sgns::crdt::Broadcaster::Next
function WaitForNext¶
Blocks until a message is queued or timeout elapses.
Parameters:
- timeout Longest time to block.
Reimplements: sgns::crdt::Broadcaster::WaitForNext
function CancelWait¶
Releases any waiter and makes later waits return at once.
Reimplements: sgns::crdt::Broadcaster::CancelWait
function Start¶
Subscribes to all configured topics and starts message processing. Must be called before using Next() to receive incoming messages.
Note: Ensures message processing is ready before any CRDT operations run.
function AddBroadcastTopic¶
Adds a new topic by name.
Parameters:
- topicName Name of the topic to add.
Return: outcome::success() on success (or if topic already existed), outcome::failure() on error.
function AddListenTopic¶
Subscribe to a given topic and store its future.
Parameters:
- topic Name of the topic to listen to.
function HasTopic¶
Checks whether the given topic is already registered.
Parameters:
- topic Name of the topic to check.
Return: True if the topic exists, false otherwise.
Reimplements: sgns::crdt::Broadcaster::HasTopic
function GetDagSyncer¶
Get the underlying GraphsyncDAGSyncer instance.
Return: Shared pointer to the GraphsyncDAGSyncer (as void pointer).
Reimplements: sgns::crdt::Broadcaster::GetDagSyncer
function Stop¶
function SetMembershipFilter¶
Installs (or replaces) the private-network membership filter consulted by OnMessage for EVERY inbound gossip message.
Parameters:
- filter Membership predicate; an empty std::function behaves like ClearMembershipFilter().
When set, a message is dropped before any CID decode, route, or queueing unless BOTH its declared protobuf peer (bmsg.peer().id()) AND its transport sender (Gossip::Message::from) pass the predicate. An empty or malformed transport from is denied under a set filter (fail-closed – mirrors sgns::networkregistry::AuthorizeGossipSender without including any networkregistry header; layering rule).
With no filter installed, OnMessage is byte-identical to the pre-filter behavior (public pass-through).
function HasMembershipFilter¶
Reports whether a membership filter is currently installed.
Return: true when OnMessage enforces membership.
function ClearMembershipFilter¶
Removes the membership filter, restoring public pass-through ingest (teardown counterpart of SetMembershipFilter).
function SetGossipSigningKey¶
Installs the gossip host keypair used to SEAL private-network publishes (CR-G01 publisher side).
Parameters:
- key The keypair that constructed the GossipPubSub host (PeerId::fromPublicKey(marshal(public key)) must equal the host's peer id, i.e. the gossip from-field it stamps).
When a membership filter is installed, Broadcast seals the serialized BroadcastMessage into an application-layer authenticated envelope (sgns::base::SealGossipPayload) signed with this keypair, and OnMessage requires every inbound message to carry a verifiable envelope whose embedded public key derives the from-field PeerId (sgns::base::OpenGossipPayload) BEFORE the membership predicate is consulted. Without a key wired, a filtered Broadcast FAILS CLOSED (publishing unsigned data that every gated receiver would deny is pointless and leaks the payload). With no filter installed, this key is unused and publish/receive stay raw and byte-identical.
function HasGossipSigningKey¶
Reports whether a gossip signing key is currently installed.
Return: true when Broadcast can seal under a set membership filter.
function AddSingleCIDInfo¶
bool AddSingleCIDInfo(
const std::string & cid,
const std::string peer_id,
const std::string address
)
function New¶
static std::shared_ptr< PubSubBroadcasterExt > New(
std::shared_ptr< sgns::crdt::GraphsyncDAGSyncer > dagSyncer,
std::shared_ptr< GossipPubSub > pubSub
)
Factory method to create a broadcaster for multiple topics.
Parameters:
- dagSyncer Graphsync DAG syncer for block exchange.
- pubSub PubSub instance used to subscribe and publish.
Return: Shared pointer to the new PubSubBroadcasterExt.
Updated on 2026-10-06 at 13:34:20 +0000