Skip to main content

EventBus

Struct EventBus 

Source
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

Source

pub fn new() -> Self

Build an empty bus.

Source

pub fn subscribe( &self, filter: EventFilter, target: SubscriptionTarget, ) -> SubscriptionId

Register a sink. Returns the id for Self::unsubscribe.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

pub fn subscribe_typed<T>(&self, tx: UnboundedSender<T>) -> SubscriptionId
where 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).

Source

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.

Source

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.

Trait Implementations§

Source§

impl Debug for EventBus

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Default for EventBus

Source§

fn default() -> EventBus

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more