mod decree_bus

module decree_bus

Transport for consensus traffic between server brains. nng transport for consensus traffic between server brains.

What travels is the capnp envelope crate::raft::wire encodes, and the bus owns that: a caller holds raft state and never a frame.

Two planes, each on the nng protocol that matches its shape. Consensus is directed, because raft addresses every message to one node, so a brain owns a Pull0 inbox that peers push into and one Push0 pipe per peer. Observation is published on a separate Pub0 socket under a topic a Sub0 reader filters by prefix, so a reader can decode the whole stream without becoming a peer of it and a slow one cannot backpressure agreement.

Receives never block. Consensus tolerates loss, reordering and duplication, which is why nothing here retries a refused send.

Structs and Unions

struct BusError(String)

Transport failure; consensus survives losses, so callers log and continue rather than unwind.

struct DecreeBus

One brain’s connection to the others over nng.

Implementations

impl DecreeBus

Functions

fn new(id: NodeId, listen_url: &str, peers: &[(NodeId, String)]) -> Result<Self, BusError>

Listen on listen_url and pipe to each (id, url) peer.

fn poll(&self) -> Vec<(NodeId, RaftMessage)>

Every message waiting for this brain, and who sent it.

fn publishing(mut self, url: &str) -> Result<Self, BusError>

Also publish a copy of everything this brain sends, on url.

A separate address from the inbox on purpose, so a reader attaches to the observation plane and never to consensus.

fn send(&self, to: NodeId, message: &RaftMessage) -> Result<(), BusError>

Deliver one message to the node it is addressed to.

struct DecreeObserver

A reader of the consensus stream that is not a peer of it.

Implementations

impl DecreeObserver

Functions

fn new(publish_urls: &[String], prefix: &str) -> Result<Self, BusError>

Watch every publisher, keeping traffic whose topic starts with prefix. The filter runs in the transport, so an unsubscribed message is never carried at all.

fn poll(&self) -> Vec<(NodeId, NodeId, RaftMessage)>

Every message seen since the last call, both ends named.