Skip to main content

lattice_plugin_host/
transient_task.rs

1//! TR.2b — the per-plugin actor bridge for transient-source providers.
2//!
3//! The transient analogue of `picker_task.rs`, and deliberately its near-twin:
4//! a dedicated async task owns the plugin's `Store<PluginState>` for life (the
5//! `Store` is `!Sync`), a [`TransientCall`] crosses an mpsc channel with a
6//! `oneshot` reply, and the `Send + Sync` [`TransientClient`] serialises calls
7//! onto the single-consumer loop.
8//!
9//! Two exports, and they are called at very different rates: `id()` once at
10//! load, to key the registry entry; `build(ctx)` once per menu open. Neither is
11//! on a hot path, which is why the design fragment can afford to call the guest
12//! per open rather than caching a spec that would go stale the moment the user
13//! moved to a different buffer.
14//!
15//! Fuel is re-armed **per call** (`arm_store` inside each `call_*`). Arming
16//! once at instantiate would be correct for a declare-once seam and wrong here:
17//! `build` is called for the life of the editor, so the menu would work for the
18//! first stretch of a session and then silently stop opening.
19//!
20//! This is the Nth near-copy of the picker / completion / decoration / media /
21//! agenda actor; the rule-of-three note in `completion_task` still stands, and
22//! this slice deliberately did not take the generalisation on mid-seam.
23
24use std::sync::Arc;
25
26use futures::StreamExt;
27use futures::channel::{mpsc, oneshot};
28use lattice_runtime::EventBus;
29use wasmtime::Store;
30
31use crate::transient_host::bindings::TransientSourcePlugin;
32use crate::{
33    Component, PluginBudget, PluginHost, PluginHostError, PluginId, PluginManifest, PluginState,
34    TrustTier, arm_store,
35};
36
37// The WIT records the bridge traffics in — the `with:`-mapped `types` mirrors,
38// i.e. the SAME Rust types `WitBoundary` round-trips. Re-exported `pub` (they
39// appear in `TransientClient`'s signatures) so callers get a clean
40// `transient_task::…` path instead of reaching into `crate::lattice::…`.
41pub use crate::lattice::plugin_host::types::{
42    TransientContext, TransientGroup, TransientItem, TransientItemKind, TransientSpec,
43};
44
45type CallResult<T> = Result<T, PluginHostError>;
46
47/// A request from a [`TransientClient`] to its [`TransientActor`].
48enum TransientCall {
49    /// `transient-source.id()` — the menu's registry name. No WIT `result`;
50    /// the reply is the name or a host-side trap.
51    Id {
52        reply: oneshot::Sender<CallResult<String>>,
53    },
54    /// `transient-source.build(ctx)` — the menu for this open. Replies the
55    /// guest's `result<transient-spec, string>` (or a host trap).
56    Build {
57        ctx: Box<TransientContext>,
58        reply: oneshot::Sender<CallResult<Result<TransientSpec, String>>>,
59    },
60}
61
62/// The `Send + Sync` handle the host adapter holds. Cloning is cheap (an mpsc
63/// `Sender` clone); every clone talks to the same actor / `Store`, so calls
64/// serialise on the single-consumer loop the `!Sync` `Store` requires.
65/// Dropping the last clone ends the actor loop — the teardown seam.
66#[derive(Clone, Debug)]
67pub struct TransientClient {
68    tx: mpsc::UnboundedSender<TransientCall>,
69    id: PluginId,
70}
71
72impl TransientClient {
73    /// The host-issued identity of the plugin behind this client.
74    pub fn id(&self) -> PluginId {
75        self.id
76    }
77
78    /// Call the guest's `id()` — the name the menu registers under. Once, at
79    /// load.
80    pub async fn menu_id(&self) -> CallResult<String> {
81        let (reply, rx) = oneshot::channel();
82        self.dispatch(TransientCall::Id { reply }, rx, "id").await
83    }
84
85    /// Call the guest's `build(ctx)`.
86    ///
87    /// The outer result is the host surface (trap / gone / quarantined); the
88    /// inner `Result<_, String>` is the guest's own WIT `result`. Both mean the
89    /// same thing to the caller — the menu does not open and the user is told
90    /// why — but they are kept distinct so the echo can say which.
91    pub async fn build(&self, ctx: TransientContext) -> CallResult<Result<TransientSpec, String>> {
92        let (reply, rx) = oneshot::channel();
93        self.dispatch(
94            TransientCall::Build {
95                ctx: Box::new(ctx),
96                reply,
97            },
98            rx,
99            "build",
100        )
101        .await
102    }
103
104    /// Shared send-then-await-reply. A closed channel (send fails) or a dropped
105    /// reply sender (the actor unwound mid-call) both surface as
106    /// [`PluginGone`](PluginHostError::PluginGone) — the caller stays live.
107    async fn dispatch<T>(
108        &self,
109        call: TransientCall,
110        rx: oneshot::Receiver<CallResult<T>>,
111        func: &'static str,
112    ) -> CallResult<T> {
113        self.tx
114            .unbounded_send(call)
115            .map_err(|_| PluginHostError::PluginGone { func })?;
116        rx.await.map_err(|_| PluginHostError::PluginGone { func })?
117    }
118}
119
120/// The per-plugin actor: owns the `Store` + transient bindings for the plugin's
121/// whole life and serves calls off the channel until every
122/// [`TransientClient`] drops.
123pub struct TransientActor {
124    store: Store<PluginState>,
125    bindings: TransientSourcePlugin,
126    budget: PluginBudget,
127    rx: mpsc::UnboundedReceiver<TransientCall>,
128    id: PluginId,
129    /// Crash-quarantine: the first export trap trips this, fires one
130    /// `PluginCrashed`, and every later call returns `Quarantined` without
131    /// re-entering the dead `Store`. A trapped builder means the menu stops
132    /// opening and says so — never a half-built menu.
133    quarantine: crate::Quarantine,
134    tracer: Option<crate::trace::PluginTracerHandle>,
135}
136
137impl TransientActor {
138    pub fn id(&self) -> PluginId {
139        self.id
140    }
141
142    /// Attach the boundary tracer (the loader calls this before spawning
143    /// `run()`). Off the hot path — the seam is async.
144    pub fn with_tracer(mut self, tracer: Option<crate::trace::PluginTracerHandle>) -> Self {
145        self.tracer = tracer;
146        self
147    }
148
149    /// Drive the actor to completion. The loop ends when the channel closes
150    /// (all clients dropped), dropping the `Store` — the teardown seam.
151    pub async fn run(mut self) {
152        while let Some(call) = self.rx.next().await {
153            match call {
154                TransientCall::Id { reply } => {
155                    let _ = reply.send(self.call_id().await);
156                }
157                TransientCall::Build { ctx, reply } => {
158                    let _ = reply.send(self.call_build(&ctx).await);
159                }
160            }
161        }
162    }
163
164    async fn call_id(&mut self) -> CallResult<String> {
165        if self.quarantine.is_tripped() {
166            return Err(PluginHostError::Quarantined { func: "id" });
167        }
168        arm_store(&mut self.store, self.budget)?;
169        let start = std::time::Instant::now();
170        let result = self
171            .bindings
172            .lattice_plugin_host_transient_source()
173            .call_id(&mut self.store)
174            .await;
175        crate::trip_and_map_traced(
176            self.tracer.as_ref(),
177            self.id.0,
178            crate::PluginSeam::TransientSource,
179            &mut self.quarantine,
180            "id",
181            start,
182            result,
183        )
184    }
185
186    async fn call_build(
187        &mut self,
188        ctx: &TransientContext,
189    ) -> CallResult<Result<TransientSpec, String>> {
190        if self.quarantine.is_tripped() {
191            return Err(PluginHostError::Quarantined { func: "build" });
192        }
193        arm_store(&mut self.store, self.budget)?;
194        let start = std::time::Instant::now();
195        let result = self
196            .bindings
197            .lattice_plugin_host_transient_source()
198            .call_build(&mut self.store, ctx)
199            .await;
200        crate::trip_and_map_traced(
201            self.tracer.as_ref(),
202            self.id.0,
203            crate::PluginSeam::TransientSource,
204            &mut self.quarantine,
205            "build",
206            start,
207            result,
208        )
209    }
210}
211
212impl PluginHost {
213    /// Instantiate a `transient-source-plugin` component under its capability
214    /// grant and return the bridge: a `Send + Sync` [`TransientClient`] plus the
215    /// [`TransientActor`] the caller drives. Grant / data-dir / WASI are
216    /// identical to every other seam (shared `build_plugin_wasi` +
217    /// `new_store`), and the actor is NOT spawned here — the lib owns no
218    /// runtime.
219    pub async fn spawn_transient_source(
220        &self,
221        component: &Component,
222        manifest: &PluginManifest,
223        tier: TrustTier,
224        budget: PluginBudget,
225        bus: &Arc<EventBus>,
226        config: Option<&Arc<lattice_config::ConfigRegistry>>,
227    ) -> Result<(TransientClient, TransientActor), PluginHostError> {
228        let (wasi, outcome, _data_dir) = self.build_plugin_wasi(manifest, tier);
229        for denied in &outcome.denied {
230            tracing::warn!(
231                plugin = %manifest.id,
232                capability = ?denied,
233                "transient plugin loaded with a withheld capability (reduced function)"
234            );
235        }
236        let mut store = self.new_store(wasi, outcome.grant, budget, Some(&manifest.id))?;
237        let bindings =
238            TransientSourcePlugin::instantiate_async(&mut store, component, &self.linker)
239                .await
240                .map_err(|e| PluginHostError::Instantiate(e.into()))?;
241        let id = self.alloc_id();
242        store.data_mut().log_ctx = self.log_ctx_for(id);
243        // A menu builder reads its OWN options — org's capture menu IS
244        // `org.capture-templates` rendered as rows. Without the registry on
245        // this store every `get-option` returns `None` and the guest falls
246        // back to its compiled default, so the menu builds from nothing and
247        // says the option is unset while `:set …?` reports a value. The
248        // `context` seam wires this for exactly the same reason; a seam that
249        // runs in its own store needs it too.
250        if let Some(registry) = config {
251            store.data_mut().config_registry = Some(Arc::clone(registry));
252        }
253        let (tx, rx) = mpsc::unbounded();
254        let client = TransientClient { tx, id };
255        let actor = TransientActor {
256            store,
257            bindings,
258            budget,
259            rx,
260            id,
261            quarantine: crate::Quarantine::new(id, Arc::clone(bus)),
262            tracer: None,
263        };
264        Ok((client, actor))
265    }
266}