Skip to main content

lattice_plugin_host/
event_task.rs

1//! PH7.8c — the per-plugin event-delivery actor + bus wiring.
2//!
3//! Design fragment: `docs/dev/architecture/plugin-host.md` §5 `events` + §3
4//! (Store-per-plugin, task-per-Store). Slice plan: PH7.8c.
5//!
6//! ## The shape
7//!
8//! Event delivery is **host→guest, fire-and-forget** — the opposite direction of
9//! the picker/completion bridges (which are host-calls-guest-for-a-reply). The
10//! host owns an mpsc; the native [`EventBus`] pushes each matched [`Event`] into
11//! it via a `SubscriptionTarget::Plugin` sink (with the bus lock dropped, so a
12//! slow handler never stalls the publisher or another subscriber, audit M1); the
13//! [`EventActor`] drains that channel on the plugin's own task and drives the
14//! guest `on-event` export. There is **no reply** — a hook observes, it does not
15//! return a value (the bus is observation-only in v1, §5.10).
16//!
17//! So this actor is simpler than [`picker_task`](crate::picker_task): no
18//! `oneshot`, no `Client` with call methods. Its input channel IS the sink the
19//! bus pushes into; the loop projects each event to WIT and calls `on-event`.
20//! Async delivery keeps plugin event work **off the keystroke path** entirely
21//! (paramount #4), bounded per-delivery by [`PluginBudget::event`].
22//!
23//! ## Runtime ownership + lifecycle
24//!
25//! The lib owns no runtime: [`PluginHost::spawn_event_plugin`] returns the
26//! `(Vec<SubscriptionId>, EventActor)` pair; the **caller** drives
27//! [`EventActor::run`] on its multi-thread runtime and holds the subscription ids
28//! to `EventBus::unsubscribe` on teardown. When every subscription is removed
29//! (teardown / lazy prune of a closed sink), the sink `Arc`s drop, the channel
30//! closes, and the actor loop ends — dropping the `Store` (the teardown seam;
31//! full crash-quarantine is PH7.12).
32
33use std::sync::Arc;
34
35use futures::StreamExt;
36use futures::channel::mpsc;
37use lattice_protocol::Event as NativeEvent;
38use lattice_runtime::{EventBus, PluginEventSink, SubscriptionId, SubscriptionTarget};
39use wasmtime::Store;
40
41use crate::WitBoundary;
42use crate::boundary_event::project_event_filter;
43use crate::events_host::bindings::EventsPlugin;
44use crate::{
45    Component, EventEmitCtx, PluginBudget, PluginHost, PluginHostError, PluginId, PluginManifest,
46    PluginState, TrustTier, arm_store, classify_trap,
47};
48
49/// One event routed from the bus to the actor. The bus sink tags each delivery
50/// with the subscription's guest-chosen `handler` id so the actor can dispatch
51/// to the right `on-event` handler (a plugin may route many `:autocmd`s to
52/// distinct handlers behind one export).
53struct PluginEventDelivery {
54    handler: u32,
55    event: NativeEvent,
56    /// OA.14d: present only for an *awaited* publish. Dropped when the guest
57    /// handler has returned — which happens on every exit path, including a
58    /// trap, a quarantine short-circuit and the actor being dropped mid-queue,
59    /// because it rides the delivery struct rather than a line of code someone
60    /// has to remember to write.
61    ack: Option<lattice_runtime::EventAck>,
62}
63
64/// A pending `wake-every` sleep, resolving to the id that came due. Boxed
65/// because the [`Sleeper`](crate::Sleeper) is a trait object — the crate owns no
66/// runtime, so the concrete future type belongs to the caller, not here.
67type BoxWake = futures::future::BoxFuture<'static, u32>;
68
69/// What the actor's `select` produced this turn.
70///
71/// OC.2 gave the loop a second input. A wake does **not** ride the bus channel:
72/// it is a future on the actor's own `FuturesUnordered`, so it needs no sender
73/// and — the load-bearing part — cannot keep the bus channel open. A plugin's
74/// actor still ends exactly when its last subscription is pruned *and* it has no
75/// wake armed, which is the property the `drop(tx)` below exists to give.
76enum Turn {
77    Event(PluginEventDelivery),
78    /// An armed `wake-every` came due. The id may since have been cancelled;
79    /// [`EventActor::deliver_wake`] is what checks.
80    Wake(u32),
81}
82
83/// The per-plugin actor: owns the `Store` + events bindings for the plugin's
84/// life and drives `on-event` for each delivery the bus pushes onto its channel.
85/// Construct via [`PluginHost::spawn_event_plugin`]; drive by spawning
86/// [`run`](Self::run) on a multi-thread runtime.
87pub struct EventActor {
88    store: Store<PluginState>,
89    bindings: EventsPlugin,
90    budget: PluginBudget,
91    rx: mpsc::UnboundedReceiver<PluginEventDelivery>,
92    id: PluginId,
93    /// Crash-quarantine (PH7.12): the first `on-event` trap trips this, firing
94    /// one `PluginCrashed` and short-circuiting every later delivery before it
95    /// re-enters the dead `Store`.
96    quarantine: crate::Quarantine,
97    /// PO.2: the boundary tracer, wired by the loader via with_tracer; None in tests / pre-wire.
98    tracer: Option<crate::trace::PluginTracerHandle>,
99}
100
101impl EventActor {
102    /// The host-issued identity of this plugin.
103    pub fn id(&self) -> PluginId {
104        self.id
105    }
106
107    /// PO.2: attach the boundary tracer (the loader calls this before spawning
108    /// run()). Off the hot path — the seam is async.
109    pub fn with_tracer(mut self, tracer: Option<crate::trace::PluginTracerHandle>) -> Self {
110        self.tracer = tracer;
111        self
112    }
113
114    /// Drive the actor to completion. Delivers each event in arrival order; the
115    /// loop ends when the channel closes (every subscription removed / pruned),
116    /// dropping the `Store`. A delivery that traps does **not** end the loop and
117    /// does **not** crash the host — the trap is caught, logged, and the loop
118    /// continues (§8). Note a component trap *taints its instance*: this plugin's
119    /// **subsequent** deliveries then also fail (each logged + skipped), so a
120    /// trapping plugin is effectively dead until it is re-instantiated
121    /// (quarantine / reload is PH7.12). The guarantee held here is **isolation**:
122    /// the publisher, the bus, every other subscriber, and every other plugin are
123    /// untouched — only the trapping plugin degrades.
124    ///
125    /// **OC.2 gave the loop a second mouth.** Besides bus deliveries it drains a
126    /// set of pending `wake-every` sleeps, firing `on-wake(id)` for each and
127    /// re-arming it. Both inputs land on this one task, so a wake is bounded by
128    /// the same budget, tripped by the same quarantine and dropped by the same
129    /// task abort as an event — none of which had to be re-implemented for it.
130    ///
131    /// The loop now ends when the channel is closed **and** nothing is armed. A
132    /// plugin that subscribes to nothing but arms a wake is a legitimate shape
133    /// (org's clock does exactly that between clock-in and clock-out), and under
134    /// the old `while let Some(..)` its actor would have exited before the first
135    /// tick.
136    pub async fn run(mut self) {
137        use futures::stream::{FusedStream, FuturesUnordered};
138
139        let mut wakes: FuturesUnordered<futures::future::BoxFuture<'static, u32>> =
140            FuturesUnordered::new();
141        // Wakes armed from inside `register-events`, before this loop existed.
142        self.arm_pending(&mut wakes);
143
144        loop {
145            if self.rx.is_terminated() && wakes.is_empty() {
146                return;
147            }
148            // Borrow `rx` explicitly so the select's futures are temporaries of
149            // this block — `deliver` below takes `&mut self`.
150            let turn = {
151                let rx = &mut self.rx;
152                if wakes.is_empty() {
153                    match rx.next().await {
154                        Some(d) => Turn::Event(d),
155                        None => continue, // re-check the exit condition above
156                    }
157                } else if rx.is_terminated() {
158                    match wakes.next().await {
159                        Some(id) => Turn::Wake(id),
160                        None => continue,
161                    }
162                } else {
163                    futures::select! {
164                        d = rx.next() => match d {
165                            Some(d) => Turn::Event(d),
166                            None => continue,
167                        },
168                        id = wakes.next() => match id {
169                            Some(id) => Turn::Wake(id),
170                            None => continue,
171                        },
172                    }
173                }
174            };
175            match turn {
176                Turn::Event(delivery) => self.deliver(delivery).await,
177                Turn::Wake(id) => self.deliver_wake(id, &mut wakes).await,
178            }
179            // A guest call may have armed more wakes (`on-event` arming one is
180            // org's clock-in path exactly), so re-check after every turn rather
181            // than only at the top.
182            self.arm_pending(&mut wakes);
183        }
184    }
185
186    /// Turn every newly-armed wake into a pending sleep. A no-op on a store with
187    /// no wake context (no `Sleeper` installed), which is why the whole seam can
188    /// be absent without the loop knowing.
189    fn arm_pending(&mut self, wakes: &mut futures::stream::FuturesUnordered<BoxWake>) {
190        let Some(ctx) = self.store.data_mut().wake.as_mut() else {
191            return;
192        };
193        let armed = ctx.take_newly_armed();
194        for id in armed {
195            if let Some(period) = ctx.period(id) {
196                wakes.push(ctx.sleep_for(id, period));
197            }
198        }
199    }
200
201    /// Fire one due wake at the guest, then re-arm it.
202    ///
203    /// Three ways this delivers nothing, each deliberate: the plugin is
204    /// quarantined (its store is dead — cancel everything so the timer stops
205    /// rather than re-entering a corpse once a minute forever); the id was
206    /// cancelled while its sleep was in flight (this is how `cancel-wake`
207    /// reaches an already-running timer); or arming the budget failed.
208    ///
209    /// Re-arming happens **after** delivery, so the interval is a gap between
210    /// wakes rather than a fixed schedule a slow guest could fall behind and
211    /// then be flooded to catch up on.
212    async fn deliver_wake(
213        &mut self,
214        id: u32,
215        wakes: &mut futures::stream::FuturesUnordered<BoxWake>,
216    ) {
217        if self.quarantine.is_tripped() {
218            if let Some(ctx) = self.store.data_mut().wake.as_mut() {
219                ctx.cancel_all();
220            }
221            return;
222        }
223        let still_armed = self
224            .store
225            .data()
226            .wake
227            .as_ref()
228            .is_some_and(|c| c.period(id).is_some());
229        if !still_armed {
230            // Cancelled mid-flight. Drop it; do not re-arm.
231            return;
232        }
233        if let Err(error) = arm_store(&mut self.store, self.budget) {
234            tracing::warn!(plugin = self.id.0, wake = id, %error, "wake skipped: arm failed");
235            return;
236        }
237        let __trace_start = std::time::Instant::now();
238        let call_result = self.bindings.call_on_wake(&mut self.store, id).await;
239        match call_result {
240            Ok(()) => {
241                if let Some(tracer) = self.tracer.as_ref() {
242                    use crate::trace::{Direction, PluginTraceRecord, TraceLevel, TraceOutcome};
243                    tracer.trace(PluginTraceRecord {
244                        plugin: self.id.0,
245                        seam: crate::PluginSeam::Events,
246                        direction: Direction::GuestExport,
247                        call: std::borrow::Cow::Borrowed("on-wake"),
248                        level: TraceLevel::Debug,
249                        outcome: TraceOutcome::Ok {
250                            micros: __trace_start.elapsed().as_micros() as u64,
251                            fuel_delta: 0,
252                        },
253                        detail: None,
254                    });
255                }
256                // Still armed? (`on-wake` itself may have cancelled it — org's
257                // clock-out does exactly that from a handler.)
258                if let Some(ctx) = self.store.data().wake.as_ref()
259                    && let Some(period) = ctx.period(id)
260                {
261                    wakes.push(ctx.sleep_for(id, period));
262                }
263            }
264            Err(source) => {
265                let kind = classify_trap(&source);
266                tracing::warn!(
267                    plugin = self.id.0,
268                    wake = id,
269                    ?kind,
270                    "plugin wake handler trapped; wake cancelled"
271                );
272                if let Some(tracer) = self.tracer.as_ref() {
273                    use crate::trace::{Direction, PluginTraceRecord, TraceLevel, TraceOutcome};
274                    tracer.trace(PluginTraceRecord {
275                        plugin: self.id.0,
276                        seam: crate::PluginSeam::Events,
277                        direction: Direction::GuestExport,
278                        call: std::borrow::Cow::Borrowed("on-wake"),
279                        level: TraceLevel::Error,
280                        outcome: TraceOutcome::Trap {
281                            kind: kind.label().to_string(),
282                            func: "on-wake".to_string(),
283                        },
284                        detail: None,
285                    });
286                }
287                self.quarantine.trip("on-wake", kind);
288                // The store is dead; a re-armed wake would only re-enter it.
289                if let Some(ctx) = self.store.data_mut().wake.as_mut() {
290                    ctx.cancel_all();
291                }
292            }
293        }
294    }
295
296    /// Deliver one event to the guest `on-event(handler, ev)` export. Every
297    /// failure mode is graceful (the four-artefact clause): a projection error
298    /// (non-UTF-8 path), an arm failure, or a guest trap (fuel/epoch/wasm) skips
299    /// *this* delivery with a `warn!`, never a panic — the plugin remains
300    /// subscribed, the publisher and every other subscriber proceed.
301    async fn deliver(&mut self, delivery: PluginEventDelivery) {
302        let PluginEventDelivery {
303            handler,
304            event,
305            // OA.14d: held for the body of this call and dropped on return —
306            // every early `return` below (quarantine, address mismatch,
307            // projection failure, arm failure) releases an awaited publisher,
308            // because "this handler will never run" is one of the ways it is
309            // over.
310            ack: _ack,
311        } = delivery;
312        // Quarantine short-circuit (PH7.12): once this instance has trapped, its
313        // `Store` is dead — skip the delivery silently (the `PluginCrashed` event
314        // already fired at trip time; re-logging every subsequent delivery is
315        // the noise this replaces).
316        if self.quarantine.is_tripped() {
317            return;
318        }
319        // OR.2: a watch batch is ADDRESSED, and this is where the address is
320        // read. The bus is a broadcast, so without this line every plugin
321        // subscribed to `files-changed` would learn which files changed under
322        // every *other* plugin's watched directory — a capability leak, since
323        // that plugin holds no `fs:read` grant over it. The id is dropped on
324        // projection, so a guest is never told its own name.
325        if let NativeEvent::FilesChanged { plugin, .. } = &event
326            && *plugin != self.id.0
327        {
328            return;
329        }
330        let wit = match event.to_wit() {
331            Ok(w) => w,
332            Err(error) => {
333                // `debug!`, not `warn!`: a wildcard-filter subscription (`kinds:
334                // none`) matches host-internal events whose `to_wit` deliberately
335                // returns Err (e.g. `Event::PluginCrashed`), so this fires once per
336                // crash — a `warn!` would flood `*messages*` per the log-levels
337                // rule. Not user-actionable: the event simply can't cross to a
338                // guest. A subscriber that wanted it would filter by a real kind.
339                tracing::debug!(
340                    plugin = self.id.0,
341                    handler,
342                    %error,
343                    "event not delivered: no WIT projection (host-internal or non-UTF-8)"
344                );
345                return;
346            }
347        };
348        if let Err(error) = arm_store(&mut self.store, self.budget) {
349            tracing::warn!(plugin = self.id.0, handler, %error, "event delivery skipped: arm failed");
350            return;
351        }
352        let __trace_start = std::time::Instant::now();
353        let call_result = self
354            .bindings
355            .call_on_event(&mut self.store, handler, &wit)
356            .await;
357        match call_result {
358            Ok(()) => {
359                // PO.2: record the successful guest-export crossing at Debug
360                // (dropped by the default Info gate — no per-delivery noise
361                // unless this plugin is raised to debug/trace). Off the hot
362                // path: the event seam is async, emission a cheap gated push.
363                if let Some(tracer) = self.tracer.as_ref() {
364                    use crate::trace::{Direction, PluginTraceRecord, TraceLevel, TraceOutcome};
365                    tracer.trace(PluginTraceRecord {
366                        plugin: self.id.0,
367                        seam: crate::PluginSeam::Events,
368                        direction: Direction::GuestExport,
369                        call: std::borrow::Cow::Borrowed("on-event"),
370                        level: TraceLevel::Debug,
371                        outcome: TraceOutcome::Ok {
372                            micros: __trace_start.elapsed().as_micros() as u64,
373                            fuel_delta: 0,
374                        },
375                        detail: None,
376                    });
377                }
378            }
379            Err(source) => {
380                // Trap (fuel/epoch/wasm) or guest panic: skip this delivery,
381                // never propagate — the host, bus, and every other subscriber
382                // are untouched (§8 isolation). The trap taints the instance
383                // irrecoverably, so trip quarantine: `PluginCrashed` fires once
384                // and every later delivery short-circuits above rather than
385                // re-failing. Re-instantiation is PH7.12b.
386                let kind = classify_trap(&source);
387                // `elapsed_ms` because an epoch trap without it says only "too
388                // long" — and the two actionable answers (the budget is wrong
389                // vs. the guest is spinning) are hundreds of milliseconds apart.
390                // Diagnosing one cost a reproduction under synthetic load.
391                tracing::warn!(
392                    plugin = self.id.0,
393                    handler,
394                    ?kind,
395                    elapsed_ms = __trace_start.elapsed().as_millis() as u64,
396                    "plugin event handler trapped; delivery skipped"
397                );
398                // PO.2: record the trapped crossing at Error (always kept),
399                // mirroring `trip_and_map_traced`'s Trap outcome shape.
400                if let Some(tracer) = self.tracer.as_ref() {
401                    use crate::trace::{Direction, PluginTraceRecord, TraceLevel, TraceOutcome};
402                    tracer.trace(PluginTraceRecord {
403                        plugin: self.id.0,
404                        seam: crate::PluginSeam::Events,
405                        direction: Direction::GuestExport,
406                        call: std::borrow::Cow::Borrowed("on-event"),
407                        level: TraceLevel::Error,
408                        outcome: TraceOutcome::Trap {
409                            kind: kind.label().to_string(),
410                            func: "on-event".to_string(),
411                        },
412                        detail: None,
413                    });
414                }
415                self.quarantine.trip("on-event", kind);
416                // OC.2: the store is dead, so every armed wake would now be a
417                // guest call into a corpse once per period, forever. Cancel them
418                // here rather than letting `deliver_wake` discover it each time.
419                if let Some(ctx) = self.store.data_mut().wake.as_mut() {
420                    ctx.cancel_all();
421                }
422            }
423        }
424    }
425}
426
427impl PluginHost {
428    /// Instantiate an `events-plugin` component under its capability grant, run
429    /// its `register-events` export to collect subscriptions, wire each to `bus`,
430    /// and return the `(subscription ids, actor)` pair. The subscriptions are
431    /// live the moment this returns — but deliveries only *fire* once the caller
432    /// drives [`EventActor::run`] (until then they queue on the channel). Grant /
433    /// data-dir / WASI are identical to
434    /// [`instantiate_plugin`](Self::instantiate_plugin) (shared `build_plugin_wasi`
435    /// + `new_store`); the actor is *not* spawned here (the lib owns no runtime).
436    ///
437    /// The caller holds the returned [`SubscriptionId`]s to `EventBus::unsubscribe`
438    /// on teardown; doing so drops the sinks, closes the channel, and ends the
439    /// actor loop.
440    ///
441    /// `config` is the live option registry, and the events seam needs it for
442    /// the same reason `context` and `transient` do — see the wiring site below.
443    pub async fn spawn_event_plugin(
444        &self,
445        component: &Component,
446        manifest: &PluginManifest,
447        tier: TrustTier,
448        budget: PluginBudget,
449        bus: &Arc<EventBus>,
450        config: Option<&Arc<lattice_config::ConfigRegistry>>,
451    ) -> Result<(Vec<SubscriptionId>, EventActor), PluginHostError> {
452        let (wasi, outcome, _data_dir) = self.build_plugin_wasi(manifest, tier);
453        for denied in &outcome.denied {
454            tracing::warn!(
455                plugin = %manifest.id,
456                capability = ?denied,
457                "event plugin loaded with a withheld capability (reduced function)"
458            );
459        }
460        let mut store = self.new_store(wasi, outcome.grant, budget, Some(&manifest.id))?;
461        let bindings = EventsPlugin::instantiate_async(&mut store, component, &self.linker)
462            .await
463            .map_err(|e| PluginHostError::Instantiate(e.into()))?;
464        let id = self.alloc_id();
465
466        // Wire the emit context BEFORE `register-events` runs: the guest may call
467        // the imported `register-event` / `emit-event` host-services from inside
468        // `register-events` (or later, from `on-event`), and both need this
469        // plugin's identity + the bus (PH7.8b.2). `store` moves into the actor
470        // below, carrying the context for the life of the plugin.
471        store.data_mut().event_emit = Some(EventEmitCtx {
472            plugin_id: id,
473            bus: Arc::clone(bus),
474        });
475        // PO.5: route this plugin's `logging` calls into the tracer (Layer 2),
476        // also before `register-events` — a guest may narrate from there.
477        store.data_mut().log_ctx = self.log_ctx_for(id);
478        // Deferred config is the whole reason `init.rs` subscribes to
479        // `plugin-loaded`: a USER plugin's options do not EXIST until it loads,
480        // so `config.set-option("org.capture-templates", …)` has to run from
481        // `on-event`, and `docs/user/init.md` documents exactly that shape.
482        // Without the registry on THIS store the call takes the
483        // "plugin has no config registry wired" branch and warns into the log,
484        // so the user's config silently does not apply and `:set …?` reports the
485        // compiled default. `context` and `transient` wire this for the same
486        // reason; a seam that runs in its own store needs it too.
487        //
488        // The gap survived because the CI.5 chain test drives the OTHER half of
489        // the documented pattern — `modes.enable-mode`, which reaches the bus
490        // rather than the registry — so the seam looked covered end to end while
491        // its config path had never been called once.
492        if let Some(registry) = config {
493            store.data_mut().config_registry = Some(Arc::clone(registry));
494        }
495
496        // OC.2: the wake context, wired BEFORE `register-events` for the same
497        // reason `event_emit` is — a guest may arm its first wake from there.
498        // `None` when no `Sleeper` was installed, which leaves `wake-every`
499        // answering `0` rather than pretending.
500        store.data_mut().wake = self
501            .sleeper
502            .get()
503            .map(|s| crate::wake::WakeCtx::new(Arc::clone(s)));
504
505        // Drive subscription registration: the guest calls the imported
506        // `events.subscribe(filter, handler)` inside `register-events`, recording
507        // each into the Store's `event_subscriptions`.
508        // PH7.8c: hold anything the guest emits from inside `register-events`.
509        //
510        // Its subscriptions are recorded into the Store during the call and
511        // wired onto the bus only AFTER it returns, so an event published in
512        // that window reaches every subscriber except the one that asked for
513        // it. A guest kicking off its own work from registration — org's roam
514        // index rings its own batch doorbell there — then waits forever for a
515        // delivery that was dropped, with nothing logged and nothing to see but
516        // a progress counter frozen at its first value.
517        //
518        // Everything a guest can CALL from here was already wired ahead of the
519        // call (emit ctx, log ctx, config, wake). This closes the last case:
520        // what it can SEND.
521        store.data_mut().deferred_events = Some(Vec::new());
522        arm_store(&mut store, budget)?;
523        bindings
524            .call_register_events(&mut store)
525            .await
526            .map_err(|source| PluginHostError::Trap {
527                func: "register-events",
528                kind: classify_trap(&source),
529                source: source.into(),
530            })?;
531        let recorded = store.data_mut().event_subscriptions.take();
532
533        // Wire each recorded subscription to the bus. The actor drains one
534        // channel; each subscription's sink tags deliveries with its handler so
535        // the actor dispatches to the right guest handler.
536        let (tx, rx) = mpsc::unbounded();
537        let mut subscription_ids = Vec::with_capacity(recorded.len());
538        for sub in recorded {
539            let filter = match project_event_filter(sub.filter) {
540                Ok(f) => f,
541                Err(error) => {
542                    tracing::warn!(
543                        plugin = id.0,
544                        handler = sub.handler,
545                        %error,
546                        "event subscription skipped: filter projection failed"
547                    );
548                    continue;
549                }
550            };
551            let handler = sub.handler;
552            let tx = tx.clone();
553            let sink: PluginEventSink = Arc::new(move |ev: NativeEvent, ack| {
554                // `false` when the actor's receiver has closed (the plugin was
555                // torn down) → the bus prunes this subscription lazily. A
556                // rejected send drops `ack` with it, so an awaited publish is
557                // released rather than left waiting on an actor that is gone.
558                tx.unbounded_send(PluginEventDelivery {
559                    handler,
560                    event: ev,
561                    ack,
562                })
563                .is_ok()
564            });
565            let sid = bus.subscribe(
566                filter,
567                SubscriptionTarget::Plugin {
568                    plugin: id.0,
569                    handler,
570                    sink,
571                },
572            );
573            subscription_ids.push(sid);
574        }
575        // Drop the original `tx`: only the sink clones (held by the live bus
576        // subscriptions) keep the channel open, so the actor ends exactly when
577        // the last subscription is unsubscribed/pruned.
578        drop(tx);
579
580        // PH7.8c: the subscriptions are live — release what the guest emitted
581        // during registration.
582        //
583        // AFTER the wiring loop and not one line earlier: the whole point is
584        // that these events find this plugin's own sinks. Publishing in
585        // registration order preserves what the guest wrote; a guest that rang
586        // two doorbells meant them in that sequence.
587        //
588        // Closing the window (back to `None`) is what makes every later emit —
589        // from `on-event`, from a wake — publish straight through, which is the
590        // behaviour that was always correct outside this window.
591        let deferred = store.data_mut().deferred_events.take().unwrap_or_default();
592        if !deferred.is_empty() {
593            tracing::debug!(
594                plugin = id.0,
595                count = deferred.len(),
596                "releasing events emitted during register-events"
597            );
598            for (name, payload) in deferred {
599                crate::host_services::emit_plugin_event(bus, name, payload);
600            }
601        }
602
603        let actor = EventActor {
604            store,
605            bindings,
606            budget,
607            rx,
608            id,
609            quarantine: crate::Quarantine::new(id, Arc::clone(bus)),
610            tracer: None,
611        };
612        Ok((subscription_ids, actor))
613    }
614}