pub struct EventBus { /* private fields */ }Expand description
Process-shared event bus. Cheap to construct (one Mutex
holding empty maps). Cheap to clone via Arc<EventBus> –
callers wrap externally; the type itself is Send + Sync and
can be shared by reference.
Implementations§
Source§impl EventBus
impl EventBus
Sourcepub fn subscribe(
&self,
filter: EventFilter,
target: SubscriptionTarget,
) -> SubscriptionId
pub fn subscribe( &self, filter: EventFilter, target: SubscriptionTarget, ) -> SubscriptionId
Register a sink. Returns the id for Self::unsubscribe.
Sourcepub fn unsubscribe(&self, id: SubscriptionId) -> bool
pub fn unsubscribe(&self, id: SubscriptionId) -> bool
Drop a previously-registered sink. Idempotent: removing an
id that has already been unsubscribed (or was never issued)
is a no-op. Returns true if at least one entry was
removed.
Sourcepub fn publish(&self, event: Event)
pub fn publish(&self, event: Event)
Fire event to every matching subscriber. Channel sinks
receive a clone (the event has Clone); Invocation sinks
queue onto Self::drain_pending_invocations for the App
to dispatch. Closed channel senders are pruned lazily on the
publish that observes them.
Locking discipline (audit M1): we hold the inner mutex only
long enough to (a) snapshot the matching channel senders
into a small Vec, and (b) push any matched Invocation
targets onto pending_invocations (the latter must stay
under the lock or it races with
Self::drain_pending_invocations). The actual tx.send
calls happen with the lock dropped, so a slow / bounded
downstream sender can never block the bus – and concurrent
publish calls from different threads can dispatch in
parallel instead of serialising on the bus mutex. That is a
deliberate relaxation of the previous “publish is totally
ordered through the bus” guarantee; within a single
publish call ordering across this caller’s subscribers is
still preserved, which is the only guarantee callers have
ever been entitled to.
EF.1: the reserved EventFilter fields are AND-checked
during the snapshot phase (under the lock) so only matching
subscribers are snapshotted / queued. A predicate therefore
runs under the bus mutex – see EventPredicate for its
non-reentrancy contract.
Sourcepub fn publish_awaited(&self, event: Event) -> Vec<Receiver<()>>
pub fn publish_awaited(&self, event: Event) -> Vec<Receiver<()>>
Self::publish with a barrier: the returned receivers each resolve
once one matched plugin subscriber’s guest handler has returned
(OA.14d). Awaiting them all is what makes an event able to precede
something — Event::PrePluginLoaded is published this way so a handler’s
set-option lands before the export that reads it.
Only plugin sinks are acknowledged. Channel subscribers are host-side
and their receivers are drained by their own owners, and Invocation
targets do not run until the App’s next turn — neither can be waited on
here without inverting the ownership the bus deliberately does not have.
A receiver resolving with Err means the handler will never run (the
actor’s channel closed, the plugin was quarantined, the task was
aborted). That is done, not a failure: the caller waits for the
handler to be over, and “it cannot run” is one of the ways it is over.
Sourcepub fn drain_pending_invocations(&self) -> Vec<CommandInvocation>
pub fn drain_pending_invocations(&self) -> Vec<CommandInvocation>
Pull every queued [CommandInvocation] target out of the
bus. The App calls this on its tick to route events that
asked the editor to “run command X” – those run through
the document actor (which is the only writer of document
state). Returned in subscription order.
Sourcepub fn subscription_count(&self) -> usize
pub fn subscription_count(&self) -> usize
Test / introspection accessor: the number of currently registered subscriptions across every kind bucket plus the wildcard list. Each multi-kind subscription is counted once per kind it touches.
Sourcepub fn subscribe_typed<T>(&self, tx: UnboundedSender<T>) -> SubscriptionIdwhere
T: TypedEvent + Clone,
pub fn subscribe_typed<T>(&self, tx: UnboundedSender<T>) -> SubscriptionIdwhere
T: TypedEvent + Clone,
M.5.3.a: subscribe to a typed event. The bus stores a
downcast closure that forwards to tx; on publish the
closure unwraps the boxed payload to the concrete type
T and sends it. Subscribers register one channel per
event type they care about; multi-type subscribers
register multiple times.
Closed channels are pruned lazily on the next matching
publish (same shape as the legacy subscribe path).
Sourcepub fn publish_typed<T>(&self, event: T)where
T: TypedEvent,
pub fn publish_typed<T>(&self, event: T)where
T: TypedEvent,
M.5.3.a: publish a typed event. Boxes the event into
Arc<dyn Any + Send + Sync> once, then walks the
TypeId-keyed subscriber bucket. Each subscriber’s
downcast closure clones the typed value into its own
channel; closures returning false (channel closed)
get pruned lazily on the next publish hitting the same
bucket.
Sourcepub fn typed_subscription_count(&self) -> usize
pub fn typed_subscription_count(&self) -> usize
M.5.3.a: count of typed subscribers across every type-id
bucket. Mirrors Self::subscription_count for the
typed surface.