pub struct Loopback {
pub(crate) id: AdapterId,
pub(crate) engine: Arc<Engine>,
pub(crate) next_sub: AtomicU64,
}Expand description
In-process handle onto the engine.
Holds its own Engine so tests can Self::subscribe and
Self::publish without going through Adapter::run. Two
handles must not share an AdapterId: Self::shutdown drops
every subscription under that id.
Fields§
§id: AdapterId§engine: Arc<Engine>§next_sub: AtomicU64Implementations§
Source§impl Loopback
impl Loopback
Sourcepub fn new(engine: Arc<Engine>, id: impl Into<AdapterId>) -> Self
pub fn new(engine: Arc<Engine>, id: impl Into<AdapterId>) -> Self
Binds this handle to engine under id.
Does not subscribe and does not start Adapter::run. The
engine is stored here; run ignores the engine the host passes
in.
Sourcepub async fn subscribe(
&self,
route: RouteKey,
) -> Result<Receiver<Envelope>, EngineError>
pub async fn subscribe( &self, route: RouteKey, ) -> Result<Receiver<Envelope>, EngineError>
Subscribes and returns a receiver of matching Envelopes.
Delivery::sub_id is stripped; this handle assigns lb-N
ids internally. Dropping the receiver stops the forwarder, but
the engine subscription stays until Self::shutdown. Further
publishes then count as dropped.
§Errors
Returns EngineError::DuplicateSub if an internal id collides,
which does not happen on a single handle.
Sourcepub async fn publish(&self, envelope: Envelope) -> PublishOutcome
pub async fn publish(&self, envelope: Envelope) -> PublishOutcome
Publishes through the shared engine. Same fan-out and drop rules
as Engine::publish.
Sourcepub async fn shutdown(&self) -> usize
pub async fn shutdown(&self) -> usize
Drops every subscription this handle owns.
Returns how many were removed. Does not cancel Adapter::run;
the host’s cancellation token does that.
Trait Implementations§
Source§impl Adapter for Loopback
impl Adapter for Loopback
Source§fn run<'async_trait>(
self: Arc<Self>,
_engine: Arc<Engine>,
shutdown: CancellationToken,
) -> Pin<Box<dyn Future<Output = Result<(), AdapterError>> + Send + 'async_trait>>where
Self: 'async_trait,
fn run<'async_trait>(
self: Arc<Self>,
_engine: Arc<Engine>,
shutdown: CancellationToken,
) -> Pin<Box<dyn Future<Output = Result<(), AdapterError>> + Send + 'async_trait>>where
Self: 'async_trait,
Waits for shutdown, then drops this handle’s subscriptions.
There is no socket. Tests that only call Loopback::subscribe
and Loopback::publish never need this.
Unlike OWP, STOMP, and DDS, this has no session to retry and
nothing in it panics in normal operation — there is no I/O and
no protocol state, only a wait on shutdown — so it carries no
on_panic/reconnect config. That is deliberate, not a gap:
panic supervision on a function that cannot fail would be
config with nothing to affect.
§Errors
Does not fail. Ok after the token fires.