pub struct DurableBroker<B: EventBroker> { /* private fields */ }Expand description
EventBroker decorator that makes every published event durable before
it is delivered: the event is journaled with fsync-before-acknowledgement
(066 FR-006) and only then forwarded to the inner broker for live
delivery, adopting the journal-assigned cursor as the inner broker’s own
cursor for that event (spec 066 FR-007) so the two stay numerically
consistent. A write that exceeds the configured timeout rejects the event
with journal_write_timeout (067 FR-003/FR-004); if the abandoned write
completes later, the writer durably revokes it so it can never surface
through replay.
subscribe/poll normally delegate to the inner broker’s fast,
in-memory delivery. When a requested cursor is older than the inner
broker’s retention window, DurableBroker falls back to the durable
journal (spec 066 FR-005, FR-008): if the journal still retains the
cursor, a durable-mode subscription is created that reads through
JournalSource::replay_from for the rest of its lifetime (applying the
identical subject_id filter as live delivery, FR-003); if the journal
has also reclaimed that history, cursor_expired is returned with the
journal’s own oldest available cursor.
Implementations§
Source§impl<B: EventBroker> DurableBroker<B>
impl<B: EventBroker> DurableBroker<B>
Sourcepub fn new(
inner: B,
sink: impl JournalSink + 'static,
source: Arc<dyn JournalSource>,
config: DurableBrokerConfig,
audit: Arc<dyn JournalWriteAuditSink>,
) -> Self
pub fn new( inner: B, sink: impl JournalSink + 'static, source: Arc<dyn JournalSource>, config: DurableBrokerConfig, audit: Arc<dyn JournalWriteAuditSink>, ) -> Self
Wrap inner with a durable write path backed by sink, and a
journal-backed replay fallback backed by source. Production callers
should use DurableBroker::open, which guarantees sink and
source share the same underlying journal storage; this lower-level
constructor exists so tests can inject independent write- and
read-side doubles.
Sourcepub fn open(
root: &Path,
inner: B,
journal_config: JournalConfig,
broker_config: DurableBrokerConfig,
audit: Arc<dyn JournalWriteAuditSink>,
clock: Arc<dyn BrokerClock>,
) -> Result<Self, JournalError>
pub fn open( root: &Path, inner: B, journal_config: JournalConfig, broker_config: DurableBrokerConfig, audit: Arc<dyn JournalWriteAuditSink>, clock: Arc<dyn BrokerClock>, ) -> Result<Self, JournalError>
Opens (or creates and recovers) a durable journal at root and wraps
inner with a write path and journal-backed replay fallback sharing
that same storage, so cursors stay consistent across live delivery
and durable replay (spec 066 FR-007).
§Errors
Returns JournalError when the journal cannot be opened or
recovered (066 FR-009).
Trait Implementations§
Source§impl<B: EventBroker> EventBroker for DurableBroker<B>
impl<B: EventBroker> EventBroker for DurableBroker<B>
Source§fn publish(&self, event: TraverseEvent) -> Result<(), EventError>
fn publish(&self, event: TraverseEvent) -> Result<(), EventError>
Active in the catalog. Read moreSource§fn subscribe(
&self,
event_type: &str,
from_cursor: &str,
) -> Result<Subscription, EventError>
fn subscribe( &self, event_type: &str, from_cursor: &str, ) -> Result<Subscription, EventError>
Source§fn subscribe_for_subject(
&self,
event_type: &str,
from_cursor: &str,
subject_id: Option<&str>,
) -> Result<Subscription, EventError>
fn subscribe_for_subject( &self, event_type: &str, from_cursor: &str, subject_id: Option<&str>, ) -> Result<Subscription, EventError>
Source§fn poll(
&self,
subscription_id: &str,
max_events: usize,
) -> Result<SubscriptionPoll, EventError>
fn poll( &self, subscription_id: &str, max_events: usize, ) -> Result<SubscriptionPoll, EventError>
max_events. Read moreSource§fn cancel(&self, subscription_id: &str) -> Result<(), EventError>
fn cancel(&self, subscription_id: &str) -> Result<(), EventError>
Source§fn publish_with_cursor(
&self,
event: TraverseEvent,
cursor: &str,
) -> Result<(), EventError>
fn publish_with_cursor( &self, event: TraverseEvent, cursor: &str, ) -> Result<(), EventError>
DurableBroker uses this
to keep live-delivery cursors numerically consistent with the durable
journal’s cursor space (spec 066 FR-007), so a cursor issued during
live polling remains valid when later resolved through durable
replay. The default implementation ignores cursor and behaves like
Self::publish. Read moreSource§fn seed_restart_floor(&self, floor: u64)
fn seed_restart_floor(&self, floor: u64)
DurableBroker::open
calls this once at construction, seeded from the durable journal’s
latest cursor, so a freshly constructed in-memory broker correctly
defers stale-looking cursors to durable replay after a restart
instead of accepting them outright. The default implementation is a
no-op.Auto Trait Implementations§
impl<B> !Freeze for DurableBroker<B>
impl<B> !RefUnwindSafe for DurableBroker<B>
impl<B> !UnwindSafe for DurableBroker<B>
impl<B> Send for DurableBroker<B>
impl<B> Sync for DurableBroker<B>
impl<B> Unpin for DurableBroker<B>where
B: Unpin,
impl<B> UnsafeUnpin for DurableBroker<B>where
B: UnsafeUnpin,
Blanket Implementations§
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> GetSetFdFlags for T
impl<T> GetSetFdFlags for T
Source§fn get_fd_flags(&self) -> Result<FdFlags, Error>where
T: AsFilelike,
fn get_fd_flags(&self) -> Result<FdFlags, Error>where
T: AsFilelike,
self file descriptor.Source§fn new_set_fd_flags(&self, fd_flags: FdFlags) -> Result<SetFdFlags<T>, Error>where
T: AsFilelike,
fn new_set_fd_flags(&self, fd_flags: FdFlags) -> Result<SetFdFlags<T>, Error>where
T: AsFilelike,
Source§fn set_fd_flags(&mut self, set_fd_flags: SetFdFlags<T>) -> Result<(), Error>where
T: Sized + AsFilelike,
fn set_fd_flags(&mut self, set_fd_flags: SetFdFlags<T>) -> Result<(), Error>where
T: Sized + AsFilelike,
self file descriptor. Read moreSource§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> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more