lattice_runtime/events.rs
1//! In-process event bus (DESIGN.md §5.10).
2//!
3//! Vim's `autocmd` and emacs's hooks both desugar to the same
4//! primitive: subscribe a sink to a filtered stream of typed
5//! events; the bus calls the sink whenever a matching event is
6//! published. v1 ships the observation-only baseline; the
7//! Before-event veto / mutation seam (§5.10.2) layers on later.
8//!
9//! ## v1 scope
10//!
11//! - Filters by [`lattice_protocol::EventKind`] for bucketing, then
12//! AND-combines the reserved [`EventFilter`] fields (`path_glob`,
13//! `major_modes`, `predicate`) at publish time on the candidates
14//! the kind bucket already selected (EF.1). Mode activation is the
15//! caller that needed them (mode-architecture.md §7.4); the checks
16//! are per-candidate constants, so dispatch stays
17//! O(subscribers-of-kind) -- the filter fields never widen the
18//! publish scan.
19//! - Sinks are [`SubscriptionTarget::Channel`] (an `mpsc::Sender`)
20//! or [`SubscriptionTarget::Invocation`] (a `CommandInvocation`
21//! the bus runs through the document actor's dispatch when the
22//! App wires that path).
23//! - Plugin handler target (`SubscriptionTarget::Plugin` in
24//! §5.10) is omitted -- WASM hosting isn't online in v1.
25//! - Indexed dispatch: subscriptions live in a
26//! `HashMap<EventKind, Vec<Subscription>>`. Publish iterates the
27//! bucket for one kind, never the global list.
28//! - Bus is `Send + Sync`; the inner state is one `Mutex`. The
29//! publish path takes the lock only to snapshot the matching
30//! channel senders (and to queue Invocation targets onto the
31//! shared `pending_invocations`); the actual `tx.send` calls
32//! run with the lock dropped. Two `publish` calls from
33//! different threads can therefore dispatch in parallel, and a
34//! future bounded subscriber cannot stall the publisher under
35//! the bus mutex (audit M1).
36//!
37//! ## What's NOT here
38//!
39//! - **`BeforeSave` / `BeforeQuit` veto.** v1 publishes the event
40//! so observers see it; mutating the payload or aborting the
41//! transition is out of scope until the actor runs the bus
42//! inside its task and respects handler errors.
43//! - **Backpressure.** Channel sinks use unbounded mpsc. If a
44//! subscriber leaks senders the bus grows. Bounded channels +
45//! slow-consumer policy follow when LSP / plugin subscribers
46//! can actually generate the volume that needs governance.
47//! - **Per-handler fuel** (§5.10.4). Plugin / Invocation handlers
48//! will eventually run with a fuel budget; v1 calls them inline
49//! or punts to caller for Invocation targets.
50
51use std::any::TypeId;
52use std::collections::HashMap;
53use std::path::Path;
54use std::sync::Arc;
55use std::sync::Mutex;
56use std::sync::atomic::{AtomicU64, Ordering};
57
58use globset::GlobSet;
59use lattice_grammar::CommandInvocation;
60use lattice_keymap::ModeId;
61use lattice_protocol::event_registry::Event as TypedEvent;
62use lattice_protocol::{Event, EventKind};
63use tokio::sync::mpsc;
64
65/// Opaque handle returned by [`EventBus::subscribe`]. Pass to
66/// [`EventBus::unsubscribe`] to remove the subscription.
67#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
68pub struct SubscriptionId(u64);
69
70impl SubscriptionId {
71 fn next() -> Self {
72 // Process-wide monotonic. One u64 is enough for the life
73 // of any process; doubles as a deterministic order for
74 // tests that need to assert dispatch order.
75 static NEXT: AtomicU64 = AtomicU64::new(1);
76 Self(NEXT.fetch_add(1, Ordering::Relaxed))
77 }
78
79 /// Raw value -- exposed for test assertions and logging only;
80 /// callers should not rely on the value beyond uniqueness.
81 pub fn raw(self) -> u64 {
82 self.0
83 }
84}
85
86/// A caller-supplied escape-hatch predicate (the `predicate` field
87/// of [`EventFilter`]). Returns `true` to keep the event.
88///
89/// **Contract (locking discipline, audit M1).** The predicate is
90/// evaluated *under the bus mutex* during the publish snapshot
91/// phase, before the lock is dropped for delivery. It must be cheap
92/// and MUST NOT re-enter the bus (no `subscribe` / `publish` /
93/// `subscription_count` from inside it) -- doing so would deadlock
94/// on the inner mutex. Modes prefer the declarative `path_glob` /
95/// `major_modes` fields (mode-architecture.md §7.4); `predicate` is
96/// for the rare case those can't express (`init.rs` custom rules).
97pub type EventPredicate = Arc<dyn Fn(&Event) -> bool + Send + Sync>;
98
99/// Filter applied at publish time. `kinds` buckets the
100/// subscription; the remaining fields (EF.1) AND-combine on top --
101/// every `Some` field must match for the event to be delivered, a
102/// `None` field is unconstrained. Mode activation
103/// (mode-architecture.md §7.4) is the caller these reserved fields
104/// were declared for.
105#[derive(Default, Clone)]
106pub struct EventFilter {
107 /// Kinds this subscription cares about. `None` means "all
108 /// kinds" -- the wildcard. Callers should always pass
109 /// `Some(...)` to keep dispatch indexed; the wildcard is
110 /// supported for one-off debugging / introspection sinks.
111 pub kinds: Option<Vec<EventKind>>,
112 /// Restrict to events whose path matches this set (e.g.
113 /// `**/*.rs` for a Rust major-mode resolver). `None` is
114 /// unconstrained. Events with no path (e.g. `BeforeQuit`,
115 /// scratch-buffer opens) never match a `Some` glob. Build the
116 /// set via [`crate::compile_glob_set`].
117 pub path_glob: Option<GlobSet>,
118 /// Restrict to events whose buffer is in one of these major
119 /// modes -- the minor-mode allowlist of §7.4. `None` is
120 /// unconstrained. No `Event` variant carries a major mode yet
121 /// (MA.1 lands `MajorEntered { major }`); until then a `Some`
122 /// allowlist matches nothing, which is the correct semantics
123 /// for "only fire inside these majors."
124 pub major_modes: Option<Vec<ModeId>>,
125 /// Restrict to minor-mode lifecycle events naming one of these
126 /// minors — the peer of [`Self::major_modes`], and the reason it
127 /// cannot simply reuse it.
128 ///
129 /// `MinorActivated` / `MinorDeactivated` carry the MINOR's name,
130 /// so `event_major_mode` answers `None` for them and a
131 /// `major_modes`-constrained filter rejects every one. Before this
132 /// field a subscriber that wanted one specific minor had no
133 /// declarative way to say so: it subscribed unfiltered and compared
134 /// names in its handler, which for a plugin means waking its task
135 /// for every minor activation in every buffer to do nothing.
136 ///
137 /// Deliberately a SEPARATE field rather than one merged `modes`
138 /// list. There are far more minors than majors and the two ask
139 /// different questions — `major_modes` means "the buffer is
140 /// entering one of these majors" (§7.4's minor-activation
141 /// allowlist), while this means "this specific minor turned on".
142 /// A merged field would answer both at once and let a subscription
143 /// fire on a name collision between the two namespaces.
144 ///
145 /// `None` is unconstrained. Both constrained is an AND, and since
146 /// no event carries both names it matches nothing — which is the
147 /// honest reading of "a major event AND a minor event".
148 pub minor_modes: Option<Vec<ModeId>>,
149 /// Escape-hatch predicate. `None` is unconstrained. See
150 /// [`EventPredicate`] for the locking contract -- it runs under
151 /// the bus mutex and must not re-enter the bus.
152 pub predicate: Option<EventPredicate>,
153}
154
155impl std::fmt::Debug for EventFilter {
156 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
157 // `predicate` is an opaque closure; report only its
158 // presence so the rest of the filter stays inspectable.
159 f.debug_struct("EventFilter")
160 .field("kinds", &self.kinds)
161 .field("path_glob", &self.path_glob)
162 .field("major_modes", &self.major_modes)
163 .field("minor_modes", &self.minor_modes)
164 .field("predicate", &self.predicate.as_ref().map(|_| "<fn>"))
165 .finish()
166 }
167}
168
169impl EventFilter {
170 /// Convenience: subscribe to a single kind.
171 pub fn kind(k: EventKind) -> Self {
172 Self {
173 kinds: Some(vec![k]),
174 ..Default::default()
175 }
176 }
177
178 /// Convenience: subscribe to a list of kinds.
179 pub fn kinds(ks: Vec<EventKind>) -> Self {
180 Self {
181 kinds: Some(ks),
182 ..Default::default()
183 }
184 }
185
186 /// Wildcard: every event.
187 pub fn any() -> Self {
188 Self {
189 kinds: None,
190 ..Default::default()
191 }
192 }
193
194 /// Builder: AND a path-glob constraint onto this filter.
195 pub fn with_path_glob(mut self, glob: GlobSet) -> Self {
196 self.path_glob = Some(glob);
197 self
198 }
199
200 /// Builder: AND a major-mode allowlist onto this filter.
201 pub fn with_major_modes(mut self, modes: Vec<ModeId>) -> Self {
202 self.major_modes = Some(modes);
203 self
204 }
205
206 /// Builder: AND an escape-hatch predicate onto this filter. See
207 /// [`EventPredicate`] for the locking contract.
208 pub fn with_predicate(mut self, predicate: EventPredicate) -> Self {
209 self.predicate = Some(predicate);
210 self
211 }
212}
213
214/// The non-`kinds` portion of an [`EventFilter`], stored per
215/// [`Subscription`] so the publish path can AND-check it against the
216/// event after the kind bucket has already selected the candidate.
217/// `kinds` is consumed by bucketing in [`EventBus::subscribe`] and
218/// never re-checked, so it is not carried here.
219#[derive(Clone, Default)]
220struct ExtraFilter {
221 path_glob: Option<GlobSet>,
222 major_modes: Option<Vec<ModeId>>,
223 minor_modes: Option<Vec<ModeId>>,
224 predicate: Option<EventPredicate>,
225}
226
227impl std::fmt::Debug for ExtraFilter {
228 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
229 f.debug_struct("ExtraFilter")
230 .field("path_glob", &self.path_glob)
231 .field("major_modes", &self.major_modes)
232 .field("minor_modes", &self.minor_modes)
233 .field("predicate", &self.predicate.as_ref().map(|_| "<fn>"))
234 .finish()
235 }
236}
237
238impl ExtraFilter {
239 /// `true` if every constrained field matches `event`. An
240 /// all-`None` filter (the common case) is unconstrained and
241 /// returns `true` immediately. Evaluated under the bus mutex --
242 /// see [`EventPredicate`] for the predicate's non-reentrancy
243 /// contract.
244 fn matches(&self, event: &Event) -> bool {
245 if let Some(glob) = &self.path_glob {
246 match event_path(event) {
247 Some(path) if glob.is_match(path) => {}
248 // No path, or path doesn't match: a path-constrained
249 // filter rejects.
250 _ => return false,
251 }
252 }
253 if let Some(allow) = &self.major_modes {
254 match event_major_mode(event) {
255 Some(major) if allow.iter().any(|id| id.as_str() == major) => {}
256 _ => return false,
257 }
258 }
259 if let Some(allow) = &self.minor_modes {
260 match event_minor_mode(event) {
261 Some(minor) if allow.iter().any(|id| id.as_str() == minor) => {}
262 _ => return false,
263 }
264 }
265 if self.predicate.as_ref().is_some_and(|p| !p(event)) {
266 return false;
267 }
268 true
269 }
270}
271
272/// A host-owned sink that delivers a matched [`Event`] to a plugin's `on-event`
273/// handler (PH7.8). Returns `false` when the plugin's receiver has closed (its
274/// actor task ended), so the bus prunes the subscription lazily — the same
275/// closed-`Channel` / `ForwardFn` discipline the rest of this module uses.
276///
277/// **The bus stays channel-agnostic.** The plugin host builds this closure over
278/// *its own* `futures::channel` mpsc (the plugin-host lib keeps `tokio` a
279/// dev-dependency; the bus's `Channel` variant uses `tokio` mpsc), so
280/// `SubscriptionTarget` never names a plugin-host or specific-channel type and
281/// `lattice-runtime` grows no dependency on either. The sink is invoked with the
282/// bus mutex **dropped** (the audit-M1 dispatch phase), so a slow plugin handler
283/// can never stall the publisher or another subscriber.
284pub type PluginEventSink = Arc<dyn Fn(Event, Option<EventAck>) -> bool + Send + Sync>;
285
286/// The completion signal an *awaited* publish threads through a plugin sink
287/// (OA.14d). The bus mints one per matched plugin subscription and
288/// [`EventBus::publish_awaited`] returns the matching receivers; the plugin's
289/// actor drops it once the guest handler has returned.
290///
291/// It is a bare `oneshot::Sender<()>` deliberately: dropping it resolves the
292/// receiver with `Err`, so a handler that traps, a store that is quarantined,
293/// and an actor whose task was aborted all release the waiter without anyone
294/// having to remember to signal. There is no way to hold the loader open by
295/// forgetting a line.
296///
297/// `None` on every ordinary publish, which is all of them but one — the bus is
298/// fire-and-forget by design, and this exists for the single event whose whole
299/// purpose is to happen *before* something else (`Event::PrePluginLoaded`).
300pub type EventAck = futures::channel::oneshot::Sender<()>;
301
302/// What the bus does when a matching event arrives.
303#[derive(Clone)]
304pub enum SubscriptionTarget {
305 /// Push the event onto an unbounded mpsc. Closed senders are
306 /// pruned lazily on the next publish that hits this kind.
307 Channel(mpsc::UnboundedSender<Event>),
308 /// Run a [`CommandInvocation`] in response. v1 does NOT execute
309 /// it (the bus has no document handle); instead it stores the
310 /// invocation and surfaces it via [`EventBus::drain_pending_invocations`]
311 /// for the App to dispatch on its turn through the actor. This
312 /// keeps the bus loop-free and side-effect-free with respect
313 /// to document state.
314 Invocation(CommandInvocation),
315 /// Deliver the event to a WASM plugin's `on-event` handler (PH7.8, filling
316 /// the reserved slot §5.10 anticipated). The host owns the `sink` (over the
317 /// plugin's actor channel); the bus calls it with the lock dropped, so a
318 /// slow plugin handler never delays the publisher or another subscriber.
319 /// `plugin` / `handler` identify the subscription for provenance and
320 /// teardown (the host unsubscribes a quarantined plugin's ids); `sink`
321 /// carries the delivery + the closed-receiver signal (`false` → prune).
322 Plugin {
323 /// The host-issued plugin id (the `u32` inside `SourceLayer::Plugin`).
324 plugin: u32,
325 /// The guest-chosen handler id, passed back to `on-event(handler, ev)`.
326 handler: u32,
327 /// The delivery sink (host-owned; runs lock-dropped).
328 sink: PluginEventSink,
329 },
330}
331
332impl std::fmt::Debug for SubscriptionTarget {
333 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
334 // `sink` is an opaque closure; report the identifying fields and mark
335 // the sink's presence (mirrors `EventFilter`'s `predicate` handling).
336 match self {
337 SubscriptionTarget::Channel(tx) => f.debug_tuple("Channel").field(tx).finish(),
338 SubscriptionTarget::Invocation(inv) => f.debug_tuple("Invocation").field(inv).finish(),
339 SubscriptionTarget::Plugin {
340 plugin, handler, ..
341 } => f
342 .debug_struct("Plugin")
343 .field("plugin", plugin)
344 .field("handler", handler)
345 .field("sink", &"<fn>")
346 .finish(),
347 }
348 }
349}
350
351#[derive(Debug)]
352struct Subscription {
353 id: SubscriptionId,
354 target: SubscriptionTarget,
355 /// EF.1: the non-`kinds` filter fields, AND-checked against the
356 /// event at publish time (after the kind bucket selected this
357 /// subscription as a candidate).
358 extra: ExtraFilter,
359}
360
361#[derive(Default)]
362struct Inner {
363 /// Indexed by kind for O(1) bucket lookup at publish.
364 by_kind: HashMap<EventKind, Vec<Subscription>>,
365 /// Wildcard subscribers (`filter.kinds == None`). Visited on
366 /// every publish; expected to be small (debug / log sinks).
367 wildcard: Vec<Subscription>,
368 /// Invocation targets the bus matched against published events
369 /// but does not own the dispatch path for. The App calls
370 /// [`EventBus::drain_pending_invocations`] each tick and routes
371 /// these through the document actor.
372 pending_invocations: Vec<CommandInvocation>,
373 /// M.5.3.a typed-event subscriptions, keyed by Rust
374 /// `std::any::TypeId`. Each entry stores a closure that
375 /// downcasts the boxed payload and forwards to the
376 /// caller's typed channel; the closure carries enough type
377 /// information that the bus's publish path stays
378 /// `Arc<dyn Any + Send + Sync>`-typed.
379 typed_subs: HashMap<TypeId, Vec<TypedSubscription>>,
380}
381
382impl std::fmt::Debug for Inner {
383 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
384 f.debug_struct("Inner")
385 .field("by_kind", &self.by_kind)
386 .field("wildcard", &self.wildcard)
387 .field("pending_invocations", &self.pending_invocations)
388 .field("typed_sub_buckets", &self.typed_subs.len())
389 .finish()
390 }
391}
392
393/// Internal record for a typed subscription. The closure
394/// downcasts the `Arc<dyn Any + Send + Sync>` payload to the
395/// concrete event type and forwards to the subscriber's typed
396/// channel; it returns `false` if the channel has been closed
397/// so the bus can prune lazily.
398///
399/// `ForwardFn` factors out the type-erased forwarder so clippy's
400/// `type_complexity` lint stops flagging the inline shape at
401/// every use site.
402type ForwardFn = Arc<dyn Fn(&Arc<dyn std::any::Any + Send + Sync>) -> bool + Send + Sync>;
403
404struct TypedSubscription {
405 id: SubscriptionId,
406 forward: ForwardFn,
407}
408
409/// Process-shared event bus. Cheap to construct (one `Mutex`
410/// holding empty maps). Cheap to clone via `Arc<EventBus>` --
411/// callers wrap externally; the type itself is `Send + Sync` and
412/// can be shared by reference.
413#[derive(Debug, Default)]
414pub struct EventBus {
415 inner: Mutex<Inner>,
416}
417
418impl EventBus {
419 /// Build an empty bus.
420 pub fn new() -> Self {
421 Self::default()
422 }
423
424 /// Register a sink. Returns the id for [`Self::unsubscribe`].
425 pub fn subscribe(&self, filter: EventFilter, target: SubscriptionTarget) -> SubscriptionId {
426 let id = SubscriptionId::next();
427 let EventFilter {
428 kinds,
429 path_glob,
430 major_modes,
431 minor_modes,
432 predicate,
433 } = filter;
434 // EF.1: the non-`kinds` fields ride along on each
435 // Subscription so publish can AND-check them. Cloned per
436 // kind bucket for multi-kind subscriptions (Arc / small-Vec
437 // clones; the common all-`None` case is free).
438 let extra = ExtraFilter {
439 path_glob,
440 major_modes,
441 minor_modes,
442 predicate,
443 };
444 let mut inner = self.inner.lock().expect("EventBus poisoned");
445 match kinds {
446 Some(kinds) => {
447 // Multi-kind subscriptions register once per kind
448 // bucket so dispatch never has to re-check the
449 // *kind* for the bucket it already matched on (the
450 // extra fields are still checked per event).
451 for k in kinds {
452 inner.by_kind.entry(k).or_default().push(Subscription {
453 id,
454 target: target.clone(),
455 extra: extra.clone(),
456 });
457 }
458 }
459 None => inner.wildcard.push(Subscription { id, target, extra }),
460 }
461 id
462 }
463
464 /// Drop a previously-registered sink. Idempotent: removing an
465 /// id that has already been unsubscribed (or was never issued)
466 /// is a no-op. Returns `true` if at least one entry was
467 /// removed.
468 pub fn unsubscribe(&self, id: SubscriptionId) -> bool {
469 let mut inner = self.inner.lock().expect("EventBus poisoned");
470 let mut removed = false;
471 for bucket in inner.by_kind.values_mut() {
472 let before = bucket.len();
473 bucket.retain(|s| s.id != id);
474 removed |= bucket.len() != before;
475 }
476 let before_wild = inner.wildcard.len();
477 inner.wildcard.retain(|s| s.id != id);
478 removed |= inner.wildcard.len() != before_wild;
479 // M.5.5: typed-event subscriptions live in their own
480 // bucket; honour the same id-keyed unsubscribe contract.
481 for bucket in inner.typed_subs.values_mut() {
482 let before = bucket.len();
483 bucket.retain(|s| s.id != id);
484 removed |= bucket.len() != before;
485 }
486 removed
487 }
488
489 /// Fire `event` to every matching subscriber. Channel sinks
490 /// receive a clone (the event has `Clone`); Invocation sinks
491 /// queue onto [`Self::drain_pending_invocations`] for the App
492 /// to dispatch. Closed channel senders are pruned lazily on the
493 /// publish that observes them.
494 ///
495 /// Locking discipline (audit M1): we hold the inner mutex only
496 /// long enough to (a) snapshot the matching channel senders
497 /// into a small Vec, and (b) push any matched Invocation
498 /// targets onto `pending_invocations` (the latter must stay
499 /// under the lock or it races with
500 /// [`Self::drain_pending_invocations`]). The actual `tx.send`
501 /// calls happen with the lock dropped, so a slow / bounded
502 /// downstream sender can never block the bus -- and concurrent
503 /// `publish` calls from different threads can dispatch in
504 /// parallel instead of serialising on the bus mutex. That is a
505 /// deliberate relaxation of the previous "publish is totally
506 /// ordered through the bus" guarantee; within a single
507 /// `publish` call ordering across this caller's subscribers is
508 /// still preserved, which is the only guarantee callers have
509 /// ever been entitled to.
510 ///
511 /// EF.1: the reserved [`EventFilter`] fields are AND-checked
512 /// during the snapshot phase (under the lock) so only matching
513 /// subscribers are snapshotted / queued. A `predicate` therefore
514 /// runs under the bus mutex -- see [`EventPredicate`] for its
515 /// non-reentrancy contract.
516 pub fn publish(&self, event: Event) {
517 self.dispatch(event, false);
518 }
519
520 /// [`Self::publish`] with a barrier: the returned receivers each resolve
521 /// once one matched *plugin* subscriber's guest handler has returned
522 /// (OA.14d). Awaiting them all is what makes an event able to precede
523 /// something — `Event::PrePluginLoaded` is published this way so a handler's
524 /// `set-option` lands before the export that reads it.
525 ///
526 /// Only plugin sinks are acknowledged. `Channel` subscribers are host-side
527 /// and their receivers are drained by their own owners, and `Invocation`
528 /// targets do not run until the App's next turn — neither can be waited on
529 /// here without inverting the ownership the bus deliberately does not have.
530 ///
531 /// A receiver resolving with `Err` means the handler will never run (the
532 /// actor's channel closed, the plugin was quarantined, the task was
533 /// aborted). That is *done*, not a failure: the caller waits for the
534 /// handler to be over, and "it cannot run" is one of the ways it is over.
535 #[must_use = "an awaited publish that is not awaited is just a publish"]
536 pub fn publish_awaited(&self, event: Event) -> Vec<futures::channel::oneshot::Receiver<()>> {
537 self.dispatch(event, true)
538 }
539
540 /// The shared body of [`Self::publish`] / [`Self::publish_awaited`]. When
541 /// `ack` is set, each plugin sink is handed an [`EventAck`] and the matching
542 /// receiver is returned.
543 fn dispatch(&self, event: Event, ack: bool) -> Vec<futures::channel::oneshot::Receiver<()>> {
544 let kind = event.kind();
545
546 // Snapshot phase: under the lock, copy out the channel
547 // senders we'll dispatch to and queue any Invocation
548 // targets. `UnboundedSender` clones are cheap (Arc bump);
549 // bucket sizes are small (subscribers per kind, plus
550 // wildcards) so this allocation is well under the cost of
551 // even one downstream `tx.send`.
552 let (channel_targets, plugin_targets) = {
553 let mut inner = self.inner.lock().expect("EventBus poisoned");
554 let mut channel_targets: Vec<(SubscriptionId, mpsc::UnboundedSender<Event>)> =
555 Vec::new();
556 // PH7.8: plugin sinks snapshotted alongside channel senders so they
557 // dispatch with the lock dropped too — a slow plugin handler's
558 // enqueue never runs under the bus mutex.
559 let mut plugin_targets: Vec<(SubscriptionId, PluginEventSink)> = Vec::new();
560
561 // Borrow-checker note: split `inner` into independent
562 // field borrows so we can read the bucket lists
563 // (immutable) while pushing onto `pending_invocations`
564 // (mutable) in one pass.
565 let Inner {
566 by_kind,
567 wildcard,
568 pending_invocations,
569 typed_subs: _,
570 } = &mut *inner;
571
572 if let Some(bucket) = by_kind.get(&kind) {
573 snapshot_bucket(bucket, &event, &mut channel_targets, &mut plugin_targets);
574 // Invocation targets: queue under the lock so we
575 // don't race with `drain_pending_invocations`.
576 // They never touch the network / channels so the
577 // cost is purely the clone, which is acceptable.
578 queue_invocations(bucket, &event, pending_invocations);
579 }
580 snapshot_bucket(wildcard, &event, &mut channel_targets, &mut plugin_targets);
581 queue_invocations(wildcard, &event, pending_invocations);
582
583 (channel_targets, plugin_targets)
584 };
585
586 // Dispatch phase: lock dropped. Slow / bounded subscribers
587 // (none today, but see the v1 backpressure note above)
588 // cannot stall the publisher under the bus mutex.
589 let mut dead: Vec<SubscriptionId> = Vec::new();
590 for (id, tx) in channel_targets {
591 if tx.send(event.clone()).is_err() {
592 dead.push(id);
593 }
594 }
595 // PH7.8: plugin sinks run lock-dropped too. A sink returning `false`
596 // (the plugin's actor channel closed) is pruned like a dead `Channel`.
597 let mut acks = Vec::new();
598 for (id, sink) in plugin_targets {
599 let ack = ack.then(|| {
600 let (tx, rx) = futures::channel::oneshot::channel();
601 acks.push(rx);
602 tx
603 });
604 if !sink(event.clone(), ack) {
605 dead.push(id);
606 }
607 }
608
609 // Pruning phase: re-acquire the lock briefly to drop dead
610 // senders. If a subscription was already removed (e.g. the
611 // owner unsubscribed concurrently) the retain is a no-op.
612 if !dead.is_empty() {
613 let mut inner = self.inner.lock().expect("EventBus poisoned");
614 if let Some(bucket) = inner.by_kind.get_mut(&kind) {
615 bucket.retain(|s| !dead.contains(&s.id));
616 }
617 inner.wildcard.retain(|s| !dead.contains(&s.id));
618 }
619 acks
620 }
621
622 /// Pull every queued [`CommandInvocation`] target out of the
623 /// bus. The App calls this on its tick to route events that
624 /// asked the editor to "run command X" -- those run through
625 /// the document actor (which is the only writer of document
626 /// state). Returned in subscription order.
627 pub fn drain_pending_invocations(&self) -> Vec<CommandInvocation> {
628 let mut inner = self.inner.lock().expect("EventBus poisoned");
629 std::mem::take(&mut inner.pending_invocations)
630 }
631
632 /// Test / introspection accessor: the number of currently
633 /// registered subscriptions across every kind bucket plus
634 /// the wildcard list. Each multi-kind subscription is counted
635 /// once per kind it touches.
636 pub fn subscription_count(&self) -> usize {
637 let inner = self.inner.lock().expect("EventBus poisoned");
638 inner.by_kind.values().map(Vec::len).sum::<usize>() + inner.wildcard.len()
639 }
640
641 /// M.5.3.a: subscribe to a *typed* event. The bus stores a
642 /// downcast closure that forwards to `tx`; on publish the
643 /// closure unwraps the boxed payload to the concrete type
644 /// `T` and sends it. Subscribers register one channel per
645 /// event type they care about; multi-type subscribers
646 /// register multiple times.
647 ///
648 /// Closed channels are pruned lazily on the next matching
649 /// publish (same shape as the legacy `subscribe` path).
650 pub fn subscribe_typed<T>(&self, tx: mpsc::UnboundedSender<T>) -> SubscriptionId
651 where
652 T: TypedEvent + Clone,
653 {
654 let id = SubscriptionId::next();
655 let forward: ForwardFn = Arc::new(move |payload| {
656 let Some(typed) = payload.downcast_ref::<T>() else {
657 // Wrong type for this subscriber -- not an error
658 // (the bus dispatches to whichever bucket it
659 // can; downcast failure should be unreachable
660 // in practice because the bucket is keyed on
661 // TypeId).
662 return true;
663 };
664 tx.send(typed.clone()).is_ok()
665 });
666 let mut inner = self.inner.lock().expect("EventBus poisoned");
667 inner
668 .typed_subs
669 .entry(TypeId::of::<T>())
670 .or_default()
671 .push(TypedSubscription { id, forward });
672 id
673 }
674
675 /// M.5.3.a: publish a typed event. Boxes the event into
676 /// `Arc<dyn Any + Send + Sync>` once, then walks the
677 /// `TypeId`-keyed subscriber bucket. Each subscriber's
678 /// downcast closure clones the typed value into its own
679 /// channel; closures returning `false` (channel closed)
680 /// get pruned lazily on the next publish hitting the same
681 /// bucket.
682 pub fn publish_typed<T>(&self, event: T)
683 where
684 T: TypedEvent,
685 {
686 let payload: Arc<dyn std::any::Any + Send + Sync> = Arc::new(event);
687 let tid = TypeId::of::<T>();
688
689 // Snapshot phase: clone the forwarder Arcs out from
690 // under the lock. Same pattern the legacy `publish`
691 // uses for channel senders -- never call into a
692 // subscriber while holding the bus mutex.
693 let forwarders: Vec<(SubscriptionId, Arc<_>)> = {
694 let inner = self.inner.lock().expect("EventBus poisoned");
695 inner
696 .typed_subs
697 .get(&tid)
698 .map(|bucket| {
699 bucket
700 .iter()
701 .map(|s| (s.id, s.forward.clone()))
702 .collect::<Vec<_>>()
703 })
704 .unwrap_or_default()
705 };
706
707 // Dispatch phase: lock dropped. Closed channels surface
708 // via `false`; we collect them for pruning.
709 let mut dead: Vec<SubscriptionId> = Vec::new();
710 for (id, forward) in forwarders {
711 if !forward(&payload) {
712 dead.push(id);
713 }
714 }
715
716 // Pruning phase.
717 if !dead.is_empty() {
718 let mut inner = self.inner.lock().expect("EventBus poisoned");
719 if let Some(bucket) = inner.typed_subs.get_mut(&tid) {
720 bucket.retain(|s| !dead.contains(&s.id));
721 }
722 }
723 }
724
725 /// M.5.3.a: count of typed subscribers across every type-id
726 /// bucket. Mirrors [`Self::subscription_count`] for the
727 /// typed surface.
728 pub fn typed_subscription_count(&self) -> usize {
729 let inner = self.inner.lock().expect("EventBus poisoned");
730 inner.typed_subs.values().map(Vec::len).sum()
731 }
732}
733
734/// Snapshot the channel senders out of one bucket so the caller
735/// can dispatch with the bus lock dropped. Cheap clone --
736/// `UnboundedSender` is internally `Arc`-backed.
737///
738/// EF.1: a subscription whose [`ExtraFilter`] rejects `event` is
739/// skipped (its `path_glob` / `major_modes` / `predicate` didn't
740/// match). The kind bucket already matched the kind; this is the
741/// per-candidate constant the publish scan adds without widening.
742fn snapshot_bucket(
743 bucket: &[Subscription],
744 event: &Event,
745 out: &mut Vec<(SubscriptionId, mpsc::UnboundedSender<Event>)>,
746 plugin_out: &mut Vec<(SubscriptionId, PluginEventSink)>,
747) {
748 for sub in bucket {
749 if !sub.extra.matches(event) {
750 continue;
751 }
752 match &sub.target {
753 SubscriptionTarget::Channel(tx) => out.push((sub.id, tx.clone())),
754 // PH7.8: plugin sinks snapshot like channel senders (a cheap `Arc`
755 // clone) so the enqueue runs with the bus lock dropped.
756 SubscriptionTarget::Plugin { sink, .. } => plugin_out.push((sub.id, sink.clone())),
757 // Invocation targets are queued separately (`queue_invocations`).
758 SubscriptionTarget::Invocation(_) => {}
759 }
760 }
761}
762
763/// Push every Invocation target in `bucket` whose [`ExtraFilter`]
764/// matches `event` onto `pending`. Stays under the bus lock (the
765/// caller holds it) because `pending_invocations` is the same field
766/// [`EventBus::drain_pending_invocations`] empties.
767fn queue_invocations(bucket: &[Subscription], event: &Event, pending: &mut Vec<CommandInvocation>) {
768 for sub in bucket {
769 if !sub.extra.matches(event) {
770 continue;
771 }
772 if let SubscriptionTarget::Invocation(inv) = &sub.target {
773 pending.push(inv.clone());
774 }
775 }
776}
777
778/// The filesystem path an event refers to, if any. EF.1's
779/// `path_glob` filter matches against this. Events with no
780/// associated path (`BeforeQuit`, `ModalModeChanged`,
781/// `OptionChanged`, scratch-buffer opens) return `None`.
782fn event_path(event: &Event) -> Option<&Path> {
783 match event {
784 Event::DocumentOpened { path, .. } => path.as_deref(),
785 Event::DocumentChanged { path, .. } => path.as_deref(),
786 Event::BeforeSave { path, .. } => Some(path.as_path()),
787 Event::DocumentSaved { path, .. } => Some(path.as_path()),
788 Event::DocumentClosed { .. }
789 | Event::SelectionsChanged { .. }
790 | Event::ModalModeChanged { .. }
791 | Event::BeforeQuit
792 | Event::OptionChanged { .. }
793 // Lifecycle events carry a buffer + mode name, not a path;
794 // they are filtered via `major_modes`, not `path_glob`.
795 | Event::MajorEntered { .. }
796 | Event::MajorExiting { .. }
797 | Event::MinorActivated { .. }
798 | Event::MinorDeactivated { .. }
799 // Plugin events carry an opaque payload, not a path; a plugin that
800 // wants path-based routing filters inside its own handler (PH7.8b).
801 | Event::Plugin { .. }
802 // A crash event carries the plugin id, not a path (PH7.12).
803 | Event::PluginCrashed { .. }
804 // Plugin-lifecycle events carry a name + id, not a path (CI.1); a
805 // handler filters by name. `PrePluginLoaded` (OA.14d) carries the name
806 // alone and is filtered the same way.
807 | Event::PrePluginLoaded { .. }
808 | Event::PluginLoaded { .. }
809 | Event::PluginUnloaded { .. }
810 // The enablement request carries a mode name, not a path (CI.4).
811 | Event::ModeEnablementRequested { .. }
812 | Event::BufferOptionOverrideRequested { .. }
813 // MG.41g: no path, and no major mode — a background task is
814 // not buffer-scoped.
815 | Event::BackgroundTaskFinished { .. }
816 // OR.2: a watch batch carries MANY paths, and this filter answers with
817 // at most one. Reporting the first would make a `path_glob` subscription
818 // match or miss a whole batch on the basis of an arbitrary member, which
819 // is worse than not matching at all — a plugin filters the batch inside
820 // its own handler, which is where it re-reads the files anyway.
821 | Event::FilesChanged { .. } => None,
822 }
823}
824
825/// The major-mode name the event's buffer is in, if the event
826/// carries one. EF.1's `major_modes` filter matches against this.
827///
828/// MA.1: `MajorEntered` / `MajorExiting` carry the major-mode name;
829/// every other variant is major-mode-agnostic and returns `None`, so
830/// a `major_modes`-constrained filter only ever matches the two
831/// lifecycle events -- the correct semantics for a minor mode that
832/// should fire when a buffer enters / exits specific majors
833/// (mode-architecture.md §7.4).
834/// The MINOR mode name a lifecycle event names, for the `minor_modes`
835/// filter. The peer of [`event_major_mode`], and deliberately disjoint
836/// from it: no event carries both, so a filter constraining both
837/// matches nothing rather than something surprising.
838fn event_minor_mode(event: &Event) -> Option<&str> {
839 match event {
840 Event::MinorActivated { minor, .. } | Event::MinorDeactivated { minor, .. } => Some(minor),
841 // Everything else is minor-mode-agnostic. Written as an
842 // explicit catch-all rather than an enumeration because the
843 // question this answers — "does this event name a minor" — has
844 // exactly two yes cases and gains nothing from listing the
845 // dozens of noes, unlike `event_major_mode` where the
846 // enumeration documents which variants were considered.
847 _ => None,
848 }
849}
850
851fn event_major_mode(event: &Event) -> Option<&str> {
852 match event {
853 Event::MajorEntered { major, .. } | Event::MajorExiting { major, .. } => Some(major),
854 // Minor lifecycle carries the *minor* name, not the major the
855 // buffer is in, so a `major_modes` filter never matches it.
856 Event::DocumentOpened { .. }
857 | Event::DocumentClosed { .. }
858 | Event::BeforeSave { .. }
859 | Event::DocumentSaved { .. }
860 | Event::DocumentChanged { .. }
861 | Event::SelectionsChanged { .. }
862 | Event::ModalModeChanged { .. }
863 | Event::BeforeQuit
864 | Event::OptionChanged { .. }
865 | Event::MinorActivated { .. }
866 | Event::MinorDeactivated { .. }
867 // A plugin event is not tied to a buffer's major mode; a plugin that
868 // wants major-mode routing filters inside its own handler (PH7.8b).
869 | Event::Plugin { .. }
870 // A crash event is not tied to a buffer's major mode (PH7.12).
871 | Event::PluginCrashed { .. }
872 // Plugin-lifecycle events are not tied to a buffer's major mode (CI.1,
873 // OA.14d).
874 | Event::PrePluginLoaded { .. }
875 | Event::PluginLoaded { .. }
876 | Event::PluginUnloaded { .. }
877 // The enablement request is not tied to a buffer's major mode (CI.4).
878 | Event::ModeEnablementRequested { .. }
879 | Event::BufferOptionOverrideRequested { .. }
880 // MG.41g: no path, and no major mode — a background task is
881 // not buffer-scoped.
882 | Event::BackgroundTaskFinished { .. }
883 // OR.2: a watch batch is about files on disk, most of which are not
884 // open in any buffer and therefore in no major mode at all.
885 | Event::FilesChanged { .. } => None,
886 }
887}
888
889#[cfg(test)]
890mod tests {
891 #![allow(clippy::unwrap_used, clippy::panic)]
892 use super::*;
893 use lattice_grammar::CommandId;
894 use lattice_protocol::ids::DocumentId;
895 use std::path::PathBuf;
896
897 fn make_event() -> Event {
898 Event::DocumentSaved {
899 id: DocumentId::new(1),
900 path: PathBuf::from("/tmp/foo.rs"),
901 }
902 }
903
904 #[test]
905 fn channel_subscriber_receives_matching_event() {
906 let bus = EventBus::new();
907 let (tx, mut rx) = mpsc::unbounded_channel();
908 bus.subscribe(
909 EventFilter::kind(EventKind::DocumentSaved),
910 SubscriptionTarget::Channel(tx),
911 );
912
913 bus.publish(make_event());
914
915 let got = rx.try_recv().expect("event delivered");
916 assert!(matches!(got, Event::DocumentSaved { .. }));
917 }
918
919 #[test]
920 fn channel_subscriber_ignores_non_matching_kinds() {
921 let bus = EventBus::new();
922 let (tx, mut rx) = mpsc::unbounded_channel();
923 bus.subscribe(
924 EventFilter::kind(EventKind::BeforeQuit),
925 SubscriptionTarget::Channel(tx),
926 );
927
928 bus.publish(make_event()); // DocumentSaved, not BeforeQuit
929 assert!(rx.try_recv().is_err(), "no event should arrive");
930
931 bus.publish(Event::BeforeQuit);
932 assert!(matches!(rx.try_recv(), Ok(Event::BeforeQuit)));
933 }
934
935 #[test]
936 fn wildcard_subscriber_sees_every_kind() {
937 let bus = EventBus::new();
938 let (tx, mut rx) = mpsc::unbounded_channel();
939 bus.subscribe(EventFilter::any(), SubscriptionTarget::Channel(tx));
940
941 bus.publish(make_event());
942 bus.publish(Event::BeforeQuit);
943
944 assert!(rx.try_recv().is_ok());
945 assert!(rx.try_recv().is_ok());
946 }
947
948 #[test]
949 fn multi_kind_subscriber_fires_for_each_listed_kind() {
950 let bus = EventBus::new();
951 let (tx, mut rx) = mpsc::unbounded_channel();
952 bus.subscribe(
953 EventFilter::kinds(vec![EventKind::BeforeSave, EventKind::DocumentSaved]),
954 SubscriptionTarget::Channel(tx),
955 );
956
957 bus.publish(Event::BeforeSave {
958 id: DocumentId::new(1),
959 path: PathBuf::from("/tmp/foo.rs"),
960 });
961 bus.publish(make_event());
962
963 assert!(matches!(rx.try_recv(), Ok(Event::BeforeSave { .. })));
964 assert!(matches!(rx.try_recv(), Ok(Event::DocumentSaved { .. })));
965 }
966
967 #[test]
968 fn unsubscribe_stops_delivery() {
969 let bus = EventBus::new();
970 let (tx, mut rx) = mpsc::unbounded_channel();
971 let id = bus.subscribe(
972 EventFilter::kind(EventKind::DocumentSaved),
973 SubscriptionTarget::Channel(tx),
974 );
975
976 assert!(bus.unsubscribe(id));
977 bus.publish(make_event());
978 assert!(rx.try_recv().is_err());
979 // Idempotent: second unsubscribe returns false.
980 assert!(!bus.unsubscribe(id));
981 }
982
983 #[test]
984 fn closed_channel_subscriber_is_pruned() {
985 let bus = EventBus::new();
986 {
987 let (tx, _rx) = mpsc::unbounded_channel();
988 bus.subscribe(
989 EventFilter::kind(EventKind::DocumentSaved),
990 SubscriptionTarget::Channel(tx),
991 );
992 }
993 // Receiver dropped; sender should now fail. Publish twice
994 // -- the first prunes the closed sub, the second confirms
995 // count is 0.
996 assert_eq!(bus.subscription_count(), 1);
997 bus.publish(make_event());
998 assert_eq!(bus.subscription_count(), 0);
999 }
1000
1001 // PH7.8: the `SubscriptionTarget::Plugin` delivery + prune surface.
1002
1003 #[test]
1004 fn plugin_target_delivers_matching_event_lock_dropped() {
1005 let bus = EventBus::new();
1006 let seen: Arc<Mutex<Vec<Event>>> = Arc::new(Mutex::new(Vec::new()));
1007 let seen_sink = Arc::clone(&seen);
1008 let sink: crate::events::PluginEventSink = Arc::new(move |ev: Event, _ack| {
1009 seen_sink.lock().expect("poisoned").push(ev);
1010 true // receiver open
1011 });
1012 bus.subscribe(
1013 EventFilter::kind(EventKind::DocumentSaved),
1014 SubscriptionTarget::Plugin {
1015 plugin: 3,
1016 handler: 9,
1017 sink,
1018 },
1019 );
1020
1021 // Non-matching kind: not delivered.
1022 bus.publish(Event::BeforeQuit);
1023 assert!(seen.lock().unwrap().is_empty());
1024
1025 // Matching kind: delivered.
1026 bus.publish(make_event());
1027 let got = seen.lock().unwrap();
1028 assert_eq!(got.len(), 1);
1029 assert!(matches!(got[0], Event::DocumentSaved { .. }));
1030 }
1031
1032 #[test]
1033 fn plugin_target_extra_filter_applies() {
1034 // The declarative filter (path_glob) gates a plugin sink identically to
1035 // a channel target (both flow through `ExtraFilter::matches`).
1036 let bus = EventBus::new();
1037 let count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1038 let count_sink = Arc::clone(&count);
1039 let sink: crate::events::PluginEventSink = Arc::new(move |_ev, _ack| {
1040 count_sink.fetch_add(1, Ordering::SeqCst);
1041 true
1042 });
1043 bus.subscribe(
1044 EventFilter::kind(EventKind::DocumentSaved)
1045 .with_path_glob(crate::compile_glob_set(["**/*.rs"])),
1046 SubscriptionTarget::Plugin {
1047 plugin: 1,
1048 handler: 1,
1049 sink,
1050 },
1051 );
1052
1053 bus.publish(saved_with_path("notes.md")); // no match
1054 bus.publish(saved_with_path("src/lib.rs")); // match
1055 assert_eq!(count.load(Ordering::SeqCst), 1);
1056 }
1057
1058 #[test]
1059 fn plugin_target_pruned_when_sink_reports_closed() {
1060 // A sink returning `false` (the plugin's actor channel closed) is pruned
1061 // lazily on the publish that observes it — the closed-`Channel` shape.
1062 let bus = EventBus::new();
1063 let sink: crate::events::PluginEventSink =
1064 Arc::new(|_ev, _ack| false /* receiver gone */);
1065 bus.subscribe(
1066 EventFilter::kind(EventKind::DocumentSaved),
1067 SubscriptionTarget::Plugin {
1068 plugin: 2,
1069 handler: 4,
1070 sink,
1071 },
1072 );
1073 assert_eq!(bus.subscription_count(), 1);
1074 bus.publish(make_event());
1075 assert_eq!(bus.subscription_count(), 0, "closed plugin sink is pruned");
1076 }
1077
1078 #[test]
1079 fn invocation_target_queues_for_drain() {
1080 let bus = EventBus::new();
1081 let inv = CommandInvocation::of(CommandId::new(42));
1082 bus.subscribe(
1083 EventFilter::kind(EventKind::DocumentSaved),
1084 SubscriptionTarget::Invocation(inv.clone()),
1085 );
1086
1087 bus.publish(make_event());
1088 let drained = bus.drain_pending_invocations();
1089 assert_eq!(drained.len(), 1);
1090 assert_eq!(drained[0].command, inv.command);
1091
1092 // Drain is destructive -- second drain returns empty.
1093 assert!(bus.drain_pending_invocations().is_empty());
1094 }
1095
1096 #[test]
1097 fn publish_does_not_hold_lock_across_dispatch() {
1098 // Audit M1 regression test. The publish path snapshots
1099 // channel senders under the lock and dispatches with the
1100 // lock dropped. Proof: a subscriber whose channel receiver
1101 // performs work that re-enters the bus (subscribing a
1102 // *new* sink) must succeed -- under the old design that
1103 // would deadlock on the inner Mutex (the publisher held
1104 // it while calling `tx.send`, which woke the receiver
1105 // task; if the receiver tried to `bus.subscribe(...)` from
1106 // the same thread of execution it would block forever).
1107 //
1108 // We model "the subscriber synchronously reacts and asks
1109 // the bus for something" via a thread that, on receiving
1110 // an event, calls `bus.subscription_count()` (which takes
1111 // the same lock). Under the new code the publisher has
1112 // already released the lock before `tx.send` fires, so
1113 // the count call returns immediately.
1114 use std::sync::Arc;
1115 use std::sync::atomic::{AtomicBool, Ordering};
1116 use std::thread;
1117 use std::time::{Duration, Instant};
1118
1119 let bus = Arc::new(EventBus::new());
1120 let (tx, mut rx) = mpsc::unbounded_channel();
1121 bus.subscribe(
1122 EventFilter::kind(EventKind::DocumentSaved),
1123 SubscriptionTarget::Channel(tx),
1124 );
1125
1126 let bus_recv = Arc::clone(&bus);
1127 let saw_event = Arc::new(AtomicBool::new(false));
1128 let saw_event_thread = Arc::clone(&saw_event);
1129 let receiver = thread::spawn(move || {
1130 // Block waiting for the event; once it arrives, take
1131 // the bus lock via `subscription_count`. If publish
1132 // were still holding the lock this would deadlock
1133 // until the test timeout below trips.
1134 let deadline = Instant::now() + Duration::from_secs(5);
1135 while rx.try_recv().is_err() {
1136 if Instant::now() > deadline {
1137 panic!("never received event");
1138 }
1139 thread::sleep(Duration::from_millis(1));
1140 }
1141 let _ = bus_recv.subscription_count();
1142 saw_event_thread.store(true, Ordering::SeqCst);
1143 });
1144
1145 bus.publish(make_event());
1146 receiver.join().expect("receiver panicked");
1147 assert!(saw_event.load(Ordering::SeqCst));
1148 }
1149
1150 #[test]
1151 fn concurrent_publishers_do_not_deadlock() {
1152 // Two threads publishing through the bus complete
1153 // independently under the new design (the lock is dropped
1154 // before dispatch, so they overlap rather than serialise
1155 // through tx.send). The assertion is just "both finish";
1156 // a regression that re-introduced lock-during-dispatch
1157 // would still pass this test, but combined with the
1158 // previous one and the channel-pruning test the new path
1159 // is well covered.
1160 use std::sync::Arc;
1161 use std::thread;
1162
1163 let bus = Arc::new(EventBus::new());
1164 let (tx, mut rx) = mpsc::unbounded_channel();
1165 bus.subscribe(EventFilter::any(), SubscriptionTarget::Channel(tx));
1166
1167 let mut handles = Vec::new();
1168 for _ in 0..4 {
1169 let b = Arc::clone(&bus);
1170 handles.push(thread::spawn(move || {
1171 for _ in 0..50 {
1172 b.publish(make_event());
1173 }
1174 }));
1175 }
1176 for h in handles {
1177 h.join().expect("publisher panicked");
1178 }
1179
1180 let mut count = 0;
1181 while rx.try_recv().is_ok() {
1182 count += 1;
1183 }
1184 assert_eq!(count, 4 * 50);
1185 }
1186
1187 #[test]
1188 fn invocation_target_persists_across_publishes() {
1189 // Unlike a channel, an Invocation subscription stays
1190 // registered after firing -- the bus didn't deliver the
1191 // payload anywhere observable, just queued an action.
1192 let bus = EventBus::new();
1193 let inv = CommandInvocation::of(CommandId::new(7));
1194 bus.subscribe(
1195 EventFilter::kind(EventKind::DocumentSaved),
1196 SubscriptionTarget::Invocation(inv),
1197 );
1198
1199 bus.publish(make_event());
1200 bus.publish(make_event());
1201 assert_eq!(bus.drain_pending_invocations().len(), 2);
1202 assert_eq!(bus.subscription_count(), 1);
1203 }
1204
1205 // EF.1: reserved-filter-field tests (path_glob / major_modes /
1206 // predicate + AND-combination + None-is-unconstrained).
1207
1208 fn saved_with_path(path: &str) -> Event {
1209 Event::DocumentSaved {
1210 id: DocumentId::new(1),
1211 path: PathBuf::from(path),
1212 }
1213 }
1214
1215 #[test]
1216 fn path_glob_filter_delivers_only_matching_paths() {
1217 let bus = EventBus::new();
1218 let (tx, mut rx) = mpsc::unbounded_channel();
1219 bus.subscribe(
1220 EventFilter::kind(EventKind::DocumentSaved)
1221 .with_path_glob(crate::compile_glob_set(["**/*.rs"])),
1222 SubscriptionTarget::Channel(tx),
1223 );
1224
1225 // Non-matching path: filtered out.
1226 bus.publish(saved_with_path("notes/todo.md"));
1227 assert!(rx.try_recv().is_err(), "*.md must not match **/*.rs");
1228
1229 // Matching path: delivered.
1230 bus.publish(saved_with_path("src/lib.rs"));
1231 assert!(matches!(rx.try_recv(), Ok(Event::DocumentSaved { .. })));
1232 }
1233
1234 #[test]
1235 fn path_glob_filter_rejects_event_with_no_path() {
1236 // A path-constrained subscription must not fire for an event
1237 // that carries no path (BeforeQuit). "Constrained on path"
1238 // means "only path-bearing events that match."
1239 let bus = EventBus::new();
1240 let (tx, mut rx) = mpsc::unbounded_channel();
1241 bus.subscribe(
1242 EventFilter::kinds(vec![EventKind::DocumentSaved, EventKind::BeforeQuit])
1243 .with_path_glob(crate::compile_glob_set(["**/*.rs"])),
1244 SubscriptionTarget::Channel(tx),
1245 );
1246
1247 bus.publish(Event::BeforeQuit);
1248 assert!(rx.try_recv().is_err(), "pathless event can't match a glob");
1249 }
1250
1251 fn major_entered(major: &str) -> Event {
1252 Event::MajorEntered {
1253 buffer: lattice_protocol::ids::BufferId::new(1),
1254 major: major.to_string(),
1255 }
1256 }
1257
1258 #[test]
1259 fn major_modes_filter_delivers_only_allowlisted_majors() {
1260 // MA.1: MajorEntered carries the major-mode name, so a
1261 // major_modes-constrained filter fires only when the entered
1262 // major is in the allowlist.
1263 let bus = EventBus::new();
1264 let (tx, mut rx) = mpsc::unbounded_channel();
1265 bus.subscribe(
1266 EventFilter::kind(EventKind::MajorEntered)
1267 .with_major_modes(vec![ModeId::new("rust-mode")]),
1268 SubscriptionTarget::Channel(tx),
1269 );
1270
1271 // Not in the allowlist: filtered out.
1272 bus.publish(major_entered("python-mode"));
1273 assert!(rx.try_recv().is_err(), "python-mode is not allowlisted");
1274
1275 // In the allowlist: delivered.
1276 bus.publish(major_entered("rust-mode"));
1277 assert!(matches!(rx.try_recv(), Ok(Event::MajorEntered { .. })));
1278 }
1279
1280 #[test]
1281 fn major_modes_filter_rejects_event_with_no_major() {
1282 // An event that carries no major (DocumentSaved) never
1283 // matches a major-mode-constrained filter.
1284 let bus = EventBus::new();
1285 let (tx, mut rx) = mpsc::unbounded_channel();
1286 bus.subscribe(
1287 EventFilter::kinds(vec![EventKind::MajorEntered, EventKind::DocumentSaved])
1288 .with_major_modes(vec![ModeId::new("rust-mode")]),
1289 SubscriptionTarget::Channel(tx),
1290 );
1291
1292 bus.publish(saved_with_path("src/lib.rs"));
1293 assert!(
1294 rx.try_recv().is_err(),
1295 "DocumentSaved carries no major, so a major_modes filter rejects it"
1296 );
1297 }
1298
1299 #[test]
1300 fn predicate_filter_gates_delivery() {
1301 let bus = EventBus::new();
1302 let (tx, mut rx) = mpsc::unbounded_channel();
1303 let pred: EventPredicate = Arc::new(
1304 |e: &Event| matches!(e, Event::DocumentSaved { path, .. } if path.ends_with("keep.rs")),
1305 );
1306 bus.subscribe(
1307 EventFilter::kind(EventKind::DocumentSaved).with_predicate(pred),
1308 SubscriptionTarget::Channel(tx),
1309 );
1310
1311 bus.publish(saved_with_path("src/drop.rs"));
1312 assert!(rx.try_recv().is_err(), "predicate rejected this event");
1313
1314 bus.publish(saved_with_path("src/keep.rs"));
1315 assert!(matches!(rx.try_recv(), Ok(Event::DocumentSaved { .. })));
1316 }
1317
1318 #[test]
1319 fn extra_fields_and_combine() {
1320 // path_glob AND predicate: both must pass.
1321 let bus = EventBus::new();
1322 let (tx, mut rx) = mpsc::unbounded_channel();
1323 let pred: EventPredicate = Arc::new(
1324 |e: &Event| matches!(e, Event::DocumentSaved { path, .. } if path.starts_with("src/")),
1325 );
1326 bus.subscribe(
1327 EventFilter::kind(EventKind::DocumentSaved)
1328 .with_path_glob(crate::compile_glob_set(["**/*.rs"]))
1329 .with_predicate(pred),
1330 SubscriptionTarget::Channel(tx),
1331 );
1332
1333 // Matches glob (*.rs) but fails predicate (not under src/).
1334 bus.publish(saved_with_path("tests/it.rs"));
1335 assert!(rx.try_recv().is_err(), "predicate half of the AND failed");
1336
1337 // Fails glob (*.md) though it would pass predicate (src/).
1338 bus.publish(saved_with_path("src/readme.md"));
1339 assert!(rx.try_recv().is_err(), "glob half of the AND failed");
1340
1341 // Both pass.
1342 bus.publish(saved_with_path("src/lib.rs"));
1343 assert!(matches!(rx.try_recv(), Ok(Event::DocumentSaved { .. })));
1344 }
1345
1346 #[test]
1347 fn no_extra_fields_is_unconstrained() {
1348 // The common case: a kinds-only filter delivers every event
1349 // of that kind (extra fields all None short-circuit to true).
1350 let bus = EventBus::new();
1351 let (tx, mut rx) = mpsc::unbounded_channel();
1352 bus.subscribe(
1353 EventFilter::kind(EventKind::DocumentSaved),
1354 SubscriptionTarget::Channel(tx),
1355 );
1356 bus.publish(saved_with_path("any/where.xyz"));
1357 assert!(matches!(rx.try_recv(), Ok(Event::DocumentSaved { .. })));
1358 }
1359
1360 #[test]
1361 fn extra_filter_applies_to_invocation_targets_too() {
1362 // The AND-check gates Invocation targets identically to
1363 // Channel targets (both flow through ExtraFilter::matches).
1364 let bus = EventBus::new();
1365 let inv = CommandInvocation::of(CommandId::new(99));
1366 bus.subscribe(
1367 EventFilter::kind(EventKind::DocumentSaved)
1368 .with_path_glob(crate::compile_glob_set(["**/*.rs"])),
1369 SubscriptionTarget::Invocation(inv),
1370 );
1371
1372 bus.publish(saved_with_path("doc.md"));
1373 assert!(
1374 bus.drain_pending_invocations().is_empty(),
1375 "non-matching path must not queue the invocation"
1376 );
1377
1378 bus.publish(saved_with_path("src/lib.rs"));
1379 assert_eq!(bus.drain_pending_invocations().len(), 1);
1380 }
1381
1382 // M.5.3.a: typed-event surface tests. Declared at module
1383 // scope so the `register_event!` macro's linkme entry lands
1384 // in the link graph.
1385 #[derive(Debug, Clone)]
1386 struct TypedTestEvent {
1387 n: u32,
1388 }
1389
1390 lattice_protocol::register_event!(
1391 TypedTestEvent,
1392 "lattice-runtime.typed-test-event",
1393 "Test event for the EventBus typed-event API.",
1394 "lattice-runtime-tests",
1395 );
1396
1397 #[test]
1398 fn typed_publish_delivers_to_typed_subscriber() {
1399 let bus = EventBus::new();
1400 let (tx, mut rx) = mpsc::unbounded_channel::<TypedTestEvent>();
1401 bus.subscribe_typed(tx);
1402 bus.publish_typed(TypedTestEvent { n: 7 });
1403 let received = rx.try_recv().expect("typed event delivered");
1404 assert_eq!(received.n, 7);
1405 }
1406
1407 #[test]
1408 fn typed_subscriber_only_sees_matching_type() {
1409 // A subscriber for one event type doesn't receive
1410 // events of another type, even when both are typed.
1411 #[derive(Debug, Clone)]
1412 struct OtherEvent {}
1413 // Can't register OtherEvent inside a fn (linkme needs
1414 // module scope). Manual impl is enough since we only
1415 // need the trait, not the descriptor entry.
1416 impl lattice_protocol::event_registry::Event for OtherEvent {
1417 fn event_type_id(&self) -> lattice_protocol::event_registry::EventTypeId {
1418 lattice_protocol::event_registry::EventTypeId::of::<Self>("test.other")
1419 }
1420 }
1421 let bus = EventBus::new();
1422 let (tx, mut rx) = mpsc::unbounded_channel::<TypedTestEvent>();
1423 bus.subscribe_typed(tx);
1424 bus.publish_typed(OtherEvent {});
1425 assert!(rx.try_recv().is_err());
1426 }
1427
1428 #[test]
1429 fn typed_subscription_count_tracks_typed_subscribers() {
1430 let bus = EventBus::new();
1431 assert_eq!(bus.typed_subscription_count(), 0);
1432 let (tx, _rx) = mpsc::unbounded_channel::<TypedTestEvent>();
1433 bus.subscribe_typed(tx);
1434 assert_eq!(bus.typed_subscription_count(), 1);
1435 }
1436
1437 #[test]
1438 fn typed_publish_prunes_dead_channel_lazily() {
1439 let bus = EventBus::new();
1440 {
1441 let (tx, _rx) = mpsc::unbounded_channel::<TypedTestEvent>();
1442 bus.subscribe_typed(tx);
1443 // Drop _rx here; tx send will fail on next publish.
1444 }
1445 // Subscriber is registered but its channel is dead.
1446 bus.publish_typed(TypedTestEvent { n: 1 });
1447 // After the publish, the dead subscriber is pruned.
1448 assert_eq!(bus.typed_subscription_count(), 0);
1449 }
1450}