Skip to main content

Subscription

Struct Subscription 

Source
pub struct Subscription<S> { /* private fields */ }
Expand description

User-facing subscription handle.

Holds a shared reference to a Store<S> backend (one Arc clone) and exposes subscribe / subscribe_all. Cheap to construct; no Arc ever appears in user code.

§Example

use core::pin::pin;
use futures::StreamExt;
use mnesis_store::{Step, StepStreamExt, Store, Subscription};

let store = Store::new(FjallStore::builder("path").open()?);
// Items are `Step<PersistedEnvelope>`: tell catch-up from live directly.
let cursor = Subscription::new(&store).subscribe(&account_id, None)?;
let mut cursor = pin!(cursor);
while let Some(item) = cursor.next().await {
    match item? {
        Step::CaughtUp => { /* backlog drained — now live */ }
        Step::Event(env) => { /* handle the raw envelope */ }
    }
}
// …or `.events()` to drop the phase, `.decoded(codec)` to decode+keep it.

Implementations§

Source§

impl<S> Subscription<S>

Source

pub fn new(store: &Store<S>) -> Self

Construct from a Store<S> handle. One Arc::clone per call.

Source§

impl<S: RawEventStore + WakeSource> Subscription<S>

Source

pub fn subscribe<I: Id>( &self, id: &I, from: Option<Version>, ) -> Result<impl Stream<Item = Result<Step<PersistedEnvelope>, <S as RawEventStore>::Error>> + Send + use<S, I>, <S as WakeSource>::Error>
where <S as RawEventStore>::Stream: Unpin,

Open a per-stream catch-up + live-tail subscription.

from: None starts from version 1; from: Some(v) starts from the event strictly after version v. Items are Step<PersistedEnvelope>: the replay events, then exactly one Step::CaughtUp at the backlog→live boundary, then live events — the phase marker is intrinsic to a subscription (a finite read has none). Checkpoint by version() on each event.

Compose from here: .events() drops the phase for a consumer that only wants events (then the full DecodedStreamExt applies); .decoded(codec) keeps the phase and hands you typed Step<Decoded<E>>. The returned stream never returns None — it waits for new events when caught up — and is !Unpin, so pin! it before polling.

§Errors

<S as WakeSource>::Error if wake-registration fails. Read errors are surfaced as Err items in the stream.

Source

pub fn subscribe_all( &self, from: Option<<S as RawEventStore>::AllPosition>, ) -> Result<impl Stream<Item = Result<Step<(<S as RawEventStore>::AllPosition, StreamKey, PersistedEnvelope)>, <S as RawEventStore>::Error>> + Send + use<S>, <S as WakeSource>::Error>
where <S as RawEventStore>::AllStream: Unpin,

Open an all-streams ($all) catch-up + live-tail subscription in AllPosition order.

from: None starts from the first event ever appended; from: Some(p) starts from the event strictly after position p. Items are Step<(AllPosition, StreamKey, PersistedEnvelope)>: the replay events, then exactly one Step::CaughtUp, then live events. Each event carries three parts beside the box: the position to checkpoint (the consumer hands it back here or to read_all to resume; the checkpoint type is adapter-defined and must be serializable), the stream key to route on without decoding the payload, and the envelope for content.

Compose with .events() / .decoded(codec) exactly as subscribe. The returned stream never returns None and is !Unpin, so pin! it before polling.

§Errors

<S as WakeSource>::Error if wake-registration fails. Read errors are surfaced as Err items in the stream.

Trait Implementations§

Source§

impl<S> Debug for Subscription<S>

Source§

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

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl<S> Freeze for Subscription<S>

§

impl<S> RefUnwindSafe for Subscription<S>
where S: RefUnwindSafe,

§

impl<S> Send for Subscription<S>
where S: Sync + Send,

§

impl<S> Sync for Subscription<S>
where S: Sync + Send,

§

impl<S> Unpin for Subscription<S>

§

impl<S> UnsafeUnpin for Subscription<S>

§

impl<S> UnwindSafe for Subscription<S>
where S: RefUnwindSafe,

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> ArchivePointee for T

Source§

type ArchivedMetadata = ()

The archived version of the pointer metadata for this type.
Source§

fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata

Converts some archived metadata to the pointer metadata for itself.
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> LayoutRaw for T

Source§

fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>

Returns the layout of the type.
Source§

impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
where T: SharedNiching<N1, N2>, N1: Niching<T>, N2: Niching<T>,

Source§

unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool

Returns whether the given value has been niched. Read more
Source§

fn resolve_niched(out: Place<NichedOption<T, N1>>)

Writes data to out indicating that a T is niched.
Source§

impl<T> Pointee for T

Source§

type Metadata = ()

The metadata type for pointers and references to this type.
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<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