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}