1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
//! The **broker** subsystem (spec section 22).
//!
//! This module contains all the components responsible for accepting
//! published messages, routing them to the correct topic, deduplicating
//! repeated submissions, maintaining per-topic offset cursors, persisting
//! message logs and snapshots, and fanning out live deliveries to
//! connected subscribers.
//!
//! # Architecture
//!
//! The central abstraction is the [`Broker`] trait, which defines the
//! full set of topic-level operations (publish, subscribe, replay, etc.)
//! as async methods. The primary implementation is [`InMemoryBroker`],
//! a single-process broker with pluggable storage backends.
//!
//! # Key components
//!
//! - **[`broker`]** — The [`Broker`] trait itself, plus the [`PublishOutcome`] type
//! and the `serialize_frame_for_fanout` helper.
//! - **[`fanout`]** — The [`FanoutEngine`] that delivers serialized frames to all
//! active subscribers of a topic.
//! - **[`router`]** — The [`TopicRouter`] trait and its [`LocalRouter`] implementation
//! that resolves topic names to [`TopicEntry`] handles.
//! - **[`memory_broker`]** — The generic [`InMemoryBroker`] struct that wires
//! all the above components together.
//!
//! [`TopicEntry`]: crate::topic::TopicEntry
/// The core broker trait and supporting types.
pub use ;
/// Fanout engine, connection sinks, subscription management, and related types.
pub use ;
/// Single-process broker with pluggable storage backends.
pub use InMemoryBroker;
/// Topic routing layer.
pub use ;