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}