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>
impl<S> Subscription<S>
Source§impl<S: RawEventStore + WakeSource> Subscription<S>
impl<S: RawEventStore + WakeSource> Subscription<S>
Sourcepub 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>
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>
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.
Sourcepub 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>
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>
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§
Auto Trait Implementations§
impl<S> Freeze for Subscription<S>
impl<S> RefUnwindSafe for Subscription<S>where
S: RefUnwindSafe,
impl<S> Send for Subscription<S>
impl<S> Sync for Subscription<S>
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> ArchivePointee for T
impl<T> ArchivePointee for T
Source§type ArchivedMetadata = ()
type ArchivedMetadata = ()
Source§fn pointer_metadata(
_: &<T as ArchivePointee>::ArchivedMetadata,
) -> <T as Pointee>::Metadata
fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> LayoutRaw for T
impl<T> LayoutRaw for T
Source§fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>
fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>
Source§impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
Source§unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool
unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool
Source§fn resolve_niched(out: Place<NichedOption<T, N1>>)
fn resolve_niched(out: Place<NichedOption<T, N1>>)
out indicating that a T is niched.