Skip to main content

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}