TopicLogSync

Struct TopicLogSync 

Source
pub struct TopicLogSync<T, S, M, L, E> {
    pub store: S,
    pub topic_map: M,
    pub topic: T,
    pub event_tx: Sender<TopicLogSyncEvent<E>>,
    pub live_mode_rx: Option<Receiver<ToSync<Operation<E>>>>,
    pub buffer_capacity: usize,
    pub _phantom: PhantomData<L>,
}
Expand description

Protocol for synchronizing logs which are associated with a generic T topic.

The mapping of T to a set of logs is handled on the application layer using an implementation of the TopicMap trait.

After sync is complete peers optionally enter “live-mode” where concurrently received and future messages will be sent directly to the application layer and forwarded to any concurrently running sync sessions. As we may receive messages from many sync sessions concurrently, messages forwarded to a sync session in live-mode are de-duplicated in order to avoid flooding the network with redundant data.

It is assumed that the T topic has been negotiated between parties prior to initiating this sync protocol.

Fields§

§store: S§topic_map: M§topic: T§event_tx: Sender<TopicLogSyncEvent<E>>§live_mode_rx: Option<Receiver<ToSync<Operation<E>>>>§buffer_capacity: usize§_phantom: PhantomData<L>

Implementations§

Source§

impl<T, S, M, L, E> TopicLogSync<T, S, M, L, E>
where T: Eq + StdHash + Serialize + for<'a> Deserialize<'a>, S: LogStore<L, E> + OperationStore<L, E>, M: TopicMap<T, Logs<L>>, L: LogId + for<'de> Deserialize<'de> + Serialize, E: Extensions,

Source

pub fn new( topic: T, store: S, topic_map: M, live_mode_rx: Option<Receiver<ToSync<Operation<E>>>>, event_tx: Sender<TopicLogSyncEvent<E>>, ) -> Self

Returns a new sync protocol instance, configured with a store and TopicMap implementation which associates the to-be-synced logs with a given topic.

Source

pub fn new_with_capacity( topic: T, store: S, topic_map: M, live_mode_rx: Option<Receiver<ToSync<Operation<E>>>>, event_tx: Sender<TopicLogSyncEvent<E>>, buffer_capacity: usize, ) -> Self

Instantiates a sync protocol with custom buffer capacity.

Trait Implementations§

Source§

impl<T: Debug, S: Debug, M: Debug, L: Debug, E: Debug> Debug for TopicLogSync<T, S, M, L, E>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl<T, S, M, L, E> Protocol for TopicLogSync<T, S, M, L, E>
where T: Debug + Eq + StdHash + Serialize + for<'a> Deserialize<'a> + Send + 'static, S: LogStore<L, E> + OperationStore<L, E> + Send + 'static, M: TopicMap<T, Logs<L>> + Send + 'static, L: LogId + for<'de> Deserialize<'de> + Serialize + Send + 'static, E: Extensions + Send + 'static,

Source§

type Error = TopicLogSyncError

Source§

type Message = TopicLogSyncMessage<L, E>

Source§

type Output = ()

Source§

async fn run( self, sink: &mut (impl Sink<Self::Message, Error = impl Debug> + Unpin), stream: &mut (impl Stream<Item = Result<Self::Message, impl Debug>> + Unpin), ) -> Result<Self::Output, Self::Error>

Auto Trait Implementations§

§

impl<T, S, M, L, E> Freeze for TopicLogSync<T, S, M, L, E>
where S: Freeze, M: Freeze, T: Freeze,

§

impl<T, S, M, L, E> !RefUnwindSafe for TopicLogSync<T, S, M, L, E>

§

impl<T, S, M, L, E> Send for TopicLogSync<T, S, M, L, E>
where S: Send, M: Send, T: Send, L: Send, E: Send,

§

impl<T, S, M, L, E> Sync for TopicLogSync<T, S, M, L, E>
where S: Sync, M: Sync, T: Sync, L: Sync, E: Send,

§

impl<T, S, M, L, E> Unpin for TopicLogSync<T, S, M, L, E>
where S: Unpin, M: Unpin, T: Unpin, L: Unpin,

§

impl<T, S, M, L, E> !UnwindSafe for TopicLogSync<T, S, M, L, E>

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
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
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.

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
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.
Source§

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

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more