Skip to main content

LazyProbabilisticBroadcast

Struct LazyProbabilisticBroadcast 

Source
pub struct LazyProbabilisticBroadcast<P, L = FairLossLink<Recovery<P>>, G = FairLossLink<Carried<Data<P>>>>
where P: Clone + Serialize + DeserializeOwned, L: VolatileLink<Recovery<P>>, L::Scope: Clone, G: VolatileLink<Carried<Data<P>>, Scope = L::Scope>,
{ /* private fields */ }
Expand description

Gossip with recovery.

Two links, both parameters: L carries the recovery traffic and G carries the gossip. Over sessions both are session links — two instances each holding one epoch per peer, both handed every scope event, on one wire — which is how crate::uniform_reliable_broadcast already puts a broadcast and a detector together. Their scopes must agree (G::Scope = L::Scope), because a scope event reaching this layer is one event about one session and goes to both.

Implementations§

Source§

impl<P> LazyProbabilisticBroadcast<P, FairLossLink<Recovery<P>>>

Source

pub fn new( me: NodeId, peers: impl IntoIterator<Item = NodeId>, config: Config, ) -> Self

Lazy probabilistic broadcast among peers, over the fair-loss links the book names.

Source§

impl<P, L> LazyProbabilisticBroadcast<P, L>

Lazy probabilistic broadcast, over the link supplied for its recovery traffic and the book’s fair-loss link for the gossip.

Only for a recovery link that reports no boundary: the gossip beneath runs over a fair-loss link here, and the two links’ scopes have to agree. Over sessions use LazyProbabilisticBroadcast::with_links with a session link for both.

Source§

impl<P, L, G> LazyProbabilisticBroadcast<P, L, G>
where P: Clone + Serialize + DeserializeOwned, L: VolatileLink<Recovery<P>>, L::Scope: Clone, G: VolatileLink<Carried<Data<P>>, Scope = L::Scope>,

Lazy probabilistic broadcast over two links: gossip beneath the eager broadcast that disseminates data, recovery for the requests and answers that repair gaps.

The link the gossip travels over.

The link the recovery traffic travels over.

Source

pub fn next_expected_of(&self, sender: Sender) -> u64

next[s], which is one until this process has delivered anything from s.

Source

pub fn next_expected(&self, from: NodeId) -> u64

next[s] for the incarnation of from most recently heard from, or one if none has been.

Source

pub fn latest(&self, from: NodeId) -> Option<Sender>

The incarnation of from most recently admitted, if any.

Source

pub fn incarnations_of(&self, from: NodeId) -> usize

How many incarnations of from this process keeps state for.

Source

pub fn pending_count(&self) -> usize

How many messages this process is holding ahead of a gap.

Source

pub fn stored_count(&self) -> usize

How many copies this process is holding for answering requests.

Source

pub fn has_stored(&self, origin: NodeId, seq: u64) -> bool

Whether this process could answer a request for seq from the incarnation of origin most recently heard from.

Trait Implementations§

Source§

impl<P, L, G> Protocol for LazyProbabilisticBroadcast<P, L, G>
where P: Clone + Serialize + DeserializeOwned, L: VolatileLink<Recovery<P>>, L::Scope: Clone, G: VolatileLink<Carried<Data<P>>, Scope = L::Scope>,

Source§

type Meta = Infallible

Keeps nothing durably.

Source§

fn on_cmd(&mut self, Cmd::Broadcast: Cmd<P>, cx: &mut ProtoCx<'_, Self>)

upon event ⟨ pb, Broadcast | m ⟩ do lsn := lsn + 1; trigger ⟨ upb, Broadcast | [DATA, …] ⟩.

Note what is not here: no delivery to self. The eager child beneath delivers a broadcast to its own process, and that arrives back through on_upb_deliver like any other, which is what puts this process’s own messages through the same sequence check as everyone else’s.

Source§

fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>)

upon event ⟨ Timeout | s, sn ⟩ do if sn > next[s] then next[s] := sn + 1.

The gap is abandoned, and sn with it — the book skips past the message the timer was waiting on, not to it. Whatever is now deliverable is released by the standing condition, which is why the drain follows.

Source§

fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>)

Both children run over the same session, so both are told when it ends or begins.

Source§

fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>)

Name this incarnation, then start the children. Runs on every restart, which is the point.

Source§

type Cmd = Cmd<P>

Requests from the layer above.
Source§

type Ind = Ind<P>

Indications to the layer above — this protocol delivering on its guarantee.
Source§

type Msg = Wire<<G as Protocol>::Msg, <L as Protocol>::Msg>

What crosses the wire to a peer running the same protocol.
Source§

type Scope = <L as Protocol>::Scope

Scopes whose boundaries this protocol’s guarantees depend on, and which it can observe. Read more
Source§

type Note = Note

The vocabulary in which this protocol narrates its decisions. Read more
Source§

type Entry = Infallible

The durable entries this protocol appends: what accumulates.
Source§

fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>)

Handle a message received from from.
§

fn on_recovery( &mut self, _cx: &mut Cx<'_, Self::Msg, Self::Ind, Self::Note, Self::Meta, Self::Entry>, )

Resume after a crash, reading what survived. Read more

Auto Trait Implementations§

§

impl<P, L, G> Freeze for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, L: Freeze, G: Freeze,

§

impl<P, L, G> RefUnwindSafe for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, L: RefUnwindSafe, G: RefUnwindSafe, <L as Protocol>::Ind: RefUnwindSafe, P: RefUnwindSafe, <G as Protocol>::Ind: RefUnwindSafe,

§

impl<P, L, G> Send for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, L: Send, P: Send, G: Send, <L as Protocol>::Ind: Send, <G as Protocol>::Ind: Send,

§

impl<P, L, G> Sync for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, L: Sync, P: Sync, G: Sync, <L as Protocol>::Ind: Sync, <G as Protocol>::Ind: Sync,

§

impl<P, L, G> Unpin for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, L: Unpin, G: Unpin, <L as Protocol>::Ind: Unpin, <G as Protocol>::Ind: Unpin, P: Unpin,

§

impl<P, L, G> UnsafeUnpin for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, L: UnsafeUnpin, G: UnsafeUnpin,

§

impl<P, L, G> UnwindSafe for LazyProbabilisticBroadcast<P, L, G>
where <L as Protocol>::Scope: Sized, P: RefUnwindSafe + UnwindSafe, L: UnwindSafe, G: UnwindSafe, <L as Protocol>::Ind: UnwindSafe, <G as Protocol>::Ind: UnwindSafe,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V