lattice_plugin_host/picker_task.rs
1//! PH7.4c.1b — the per-plugin actor task + call protocol (the bridge).
2//!
3//! Design fragment: `docs/dev/architecture/plugin-host.md` §3 (the
4//! Store-per-plugin, task-per-Store model) and §4.3 (host owns the async
5//! plumbing). Slice plan: `slice-plans/plugin-host.md` PH7.4c.1b.
6//!
7//! ## The problem this solves
8//!
9//! A picker source's exports (`spec`/`init`/`accept`, `picker_host.rs`) are
10//! **async, single-threaded, and fuel-bounded**: they run against a
11//! `wasmtime::Store<PluginState>`, which is `!Sync` and must not be touched
12//! concurrently. But the host wraps a source as `Arc<dyn PickerSourceGenerator>`
13//! (PH7.4c.2) — a `Send + Sync` trait object the picker registry calls from
14//! anywhere. Something has to bridge a `Send + Sync` caller to a `!Sync`,
15//! single-threaded, async guest.
16//!
17//! ## The shape (locked with Dhruva, design option A)
18//!
19//! Each plugin runs as **one dedicated async task** ([`PickerActor`]) that owns
20//! its `Store` for its whole life. The task loops over an mpsc channel; each
21//! [`PickerCall`] carries the call's inputs plus a `oneshot` reply sender. The
22//! `Send + Sync` [`PickerClient`] holds the mpsc `Sender`: `init`/`accept`/`spec`
23//! send a request and await the oneshot reply. The `Store` is never locked and
24//! never leaves the task, so it stays single-threaded by construction; per-call
25//! fuel/epoch is armed inside the loop, before each guest call.
26//!
27//! This is chosen over `Arc<async_mutex<Store>>` (a lock held across `.await`,
28//! and a runtime dependency baked into the lib) because it keeps the `Store`
29//! genuinely single-threaded and reuses cleanly for every future guest-backed
30//! `Arc<dyn>` adapter (completion, grammar, modes). See the fragment's
31//! heuristic-#1 note.
32//!
33//! ## Runtime ownership
34//!
35//! The lib owns no async runtime. [`PluginHost::spawn_picker_source`] returns
36//! the `(PickerClient, PickerActor)` pair; the **caller** drives the actor by
37//! spawning [`PickerActor::run`] on its own multi-thread runtime (never the
38//! `current_thread` editor actor — paramount #4 + the no-UI-thread-work rule).
39//! The channels are `futures::channel`, runtime-agnostic on purpose.
40
41use std::sync::Arc;
42
43use futures::StreamExt;
44use futures::channel::{mpsc, oneshot};
45use lattice_runtime::EventBus;
46use wasmtime::Store;
47
48use crate::picker_host::bindings::PickerSourcePlugin;
49use crate::{
50 Component, PluginBudget, PluginHost, PluginHostError, PluginId, PluginManifest, PluginState,
51 TrustTier, arm_store,
52};
53/// OR.5b: the NATIVE spec. `register-picker-source` converts each spec at the
54/// host-import call, so the actor hands back specs that have already crossed —
55/// converting again here would be doing the same work twice and giving the
56/// second attempt a chance to disagree.
57use lattice_picker::source::PickerSourceSpec as NativePickerSourceSpec;
58
59// The picker WIT records the bridge's public API traffics in. They are the
60// `with:`-mapped `types` mirrors (`picker_host.rs`) plus the picker-interface
61// `candidate-pair`, i.e. the SAME Rust types `WitBoundary` round-trips — the
62// boundary conversion (native ↔ these) is PH7.4c.2's job, so this bridge speaks
63// purely in WIT types and stays conversion-free. Re-exported `pub` (they appear
64// in the `PickerClient` method signatures) so callers get a clean
65// `picker_task::…` path instead of reaching into `crate::lattice::…`.
66pub use crate::lattice::plugin_host::types::{
67 ActiveBufferSnapshot, PickerAcceptOutcome, PickerContext, PickerSourceSpec, Position,
68 RoutingPayload,
69};
70pub use crate::picker_host::bindings::exports::lattice::plugin_host::picker_source::CandidatePair;
71
72/// The result of a picker guest call routed through the actor. The outer
73/// [`PluginHostError`] is the *host-side* failure surface — a
74/// [`Trap`](PluginHostError::Trap) (fuel/epoch/wasm) or
75/// [`PluginGone`](PluginHostError::PluginGone) (the actor ended). The inner
76/// `Result<T, String>` is the guest's own typed WIT `result` — an `init` that
77/// declined, an `accept` that could not resolve its routing token. `spec` has
78/// no WIT `result`, so its call type is just `Result<PickerSourceSpec, _>`.
79type CallResult<T> = Result<T, PluginHostError>;
80
81/// A request sent from a [`PickerClient`] to its [`PickerActor`]. Each variant
82/// carries the guest inputs plus the `oneshot` the actor replies on. The large
83/// [`PickerContext`] projection is boxed so the enum stays small.
84enum PickerCall {
85 /// OR.5b: `register-picker-sources()` — drive the guest's registration
86 /// export, then hand back every spec it declared through the imported
87 /// `register-picker-source`.
88 ///
89 /// Replaces `spec()`. The difference is the slice: a component used to BE
90 /// one source and now DECLARES N, so registration is a call the guest makes
91 /// rather than a value the host reads.
92 RegisterSources {
93 reply: oneshot::Sender<CallResult<Vec<NativePickerSourceSpec>>>,
94 },
95 /// `picker-source.init(source, ctx, args)` — build the candidate set for
96 /// ONE of this plugin's sources. Replies the guest's
97 /// `result<list<candidate-pair>, string>` (or a host trap).
98 Init {
99 source: String,
100 ctx: Box<PickerContext>,
101 args: Vec<String>,
102 reply: oneshot::Sender<CallResult<Result<Vec<CandidatePair>, String>>>,
103 },
104 /// `picker-source.accept(source, ctx, routing)` — resolve a chosen routing
105 /// token. Replies the guest's `result<picker-accept-outcome, string>`.
106 Accept {
107 source: String,
108 ctx: Box<PickerContext>,
109 routing: RoutingPayload,
110 reply: oneshot::Sender<CallResult<Result<PickerAcceptOutcome, String>>>,
111 },
112}
113
114/// The `Send + Sync` handle a host adapter (PH7.4c.2) holds. Cloning it is cheap
115/// (an mpsc `Sender` clone); every clone talks to the same actor / `Store`, so
116/// calls are serialized by the single-consumer loop — the guarantee the `!Sync`
117/// `Store` needs. Dropping the last clone ends the actor loop (teardown).
118#[derive(Clone)]
119pub struct PickerClient {
120 tx: mpsc::UnboundedSender<PickerCall>,
121 id: PluginId,
122}
123
124impl PickerClient {
125 /// The host-issued identity of the plugin behind this client — the `u32`
126 /// inside its `SourceLayer::Plugin(id)` provenance.
127 pub fn id(&self) -> PluginId {
128 self.id
129 }
130
131 /// OR.5b: drive the guest's `register-picker-sources()` and collect every
132 /// source it declared. Returns a typed host error
133 /// ([`PluginGone`](PluginHostError::PluginGone) if the actor has ended,
134 /// [`Trap`](PluginHostError::Trap) on fuel/epoch/wasm).
135 ///
136 /// An empty list is not an error: a plugin that provides the seam and
137 /// declares nothing registers nothing, which is what it asked for.
138 pub async fn register_sources(&self) -> CallResult<Vec<NativePickerSourceSpec>> {
139 let (reply, rx) = oneshot::channel();
140 self.dispatch(
141 PickerCall::RegisterSources { reply },
142 rx,
143 "register-picker-sources",
144 )
145 .await
146 }
147
148 /// Call the guest's `init(source, ctx, args)`. The outer result is the host
149 /// surface; the inner `Result<_, String>` is the guest's own WIT `result`
150 /// (an `Err` string is a source that declined to produce candidates).
151 pub async fn init(
152 &self,
153 source: String,
154 ctx: PickerContext,
155 args: Vec<String>,
156 ) -> CallResult<Result<Vec<CandidatePair>, String>> {
157 let (reply, rx) = oneshot::channel();
158 self.dispatch(
159 PickerCall::Init {
160 source,
161 ctx: Box::new(ctx),
162 args,
163 reply,
164 },
165 rx,
166 "init",
167 )
168 .await
169 }
170
171 /// Call the guest's `accept(source, ctx, routing)` — translate the user's
172 /// chosen routing token into a typed outcome the host applies.
173 pub async fn accept(
174 &self,
175 source: String,
176 ctx: PickerContext,
177 routing: RoutingPayload,
178 ) -> CallResult<Result<PickerAcceptOutcome, String>> {
179 let (reply, rx) = oneshot::channel();
180 self.dispatch(
181 PickerCall::Accept {
182 source,
183 ctx: Box::new(ctx),
184 routing,
185 reply,
186 },
187 rx,
188 "accept",
189 )
190 .await
191 }
192
193 /// Shared send-then-await-reply. A closed channel (send fails) or a dropped
194 /// reply sender (the actor unwound mid-call) both surface as
195 /// [`PluginGone`](PluginHostError::PluginGone) — the caller stays live.
196 async fn dispatch<T>(
197 &self,
198 call: PickerCall,
199 rx: oneshot::Receiver<CallResult<T>>,
200 func: &'static str,
201 ) -> CallResult<T> {
202 self.tx
203 .unbounded_send(call)
204 .map_err(|_| PluginHostError::PluginGone { func })?;
205 rx.await.map_err(|_| PluginHostError::PluginGone { func })?
206 }
207}
208
209/// The per-plugin actor: owns the `Store` + picker bindings for the plugin's
210/// whole life and serves calls off the channel until every [`PickerClient`] is
211/// dropped. Construct via [`PluginHost::spawn_picker_source`]; drive by spawning
212/// [`run`](Self::run) on a multi-thread runtime.
213pub struct PickerActor {
214 store: Store<PluginState>,
215 bindings: PickerSourcePlugin,
216 budget: PluginBudget,
217 rx: mpsc::UnboundedReceiver<PickerCall>,
218 id: PluginId,
219 /// Crash-quarantine (PH7.12): the first export trap trips this, fires one
220 /// `PluginCrashed`, and every later call returns `Quarantined` without
221 /// re-entering the dead `Store`.
222 quarantine: crate::Quarantine,
223 /// PO.2: the boundary tracer, wired by the loader via with_tracer; None in tests / pre-wire.
224 tracer: Option<crate::trace::PluginTracerHandle>,
225}
226
227impl PickerActor {
228 /// The host-issued identity of this plugin (matches its [`PickerClient::id`]).
229 pub fn id(&self) -> PluginId {
230 self.id
231 }
232
233 /// PO.2: attach the boundary tracer (the loader calls this before spawning
234 /// run()). Off the hot path — the seam is async.
235 pub fn with_tracer(mut self, tracer: Option<crate::trace::PluginTracerHandle>) -> Self {
236 self.tracer = tracer;
237 self
238 }
239
240 /// Drive the actor to completion. Serves each [`PickerCall`] in arrival
241 /// order — arming the per-call fuel/epoch budget, calling the guest export,
242 /// mapping a trap to a typed error — and replies on the call's `oneshot`.
243 /// A trap does **not** end the loop: the `Store` survives a clean fuel/epoch
244 /// trap, so the source stays callable (full crash-quarantine is PH7.12). The
245 /// loop ends when the channel closes (all clients dropped), dropping the
246 /// `Store` — the teardown seam. If a caller has already dropped its reply
247 /// receiver, the send is a no-op (the call was abandoned).
248 pub async fn run(mut self) {
249 while let Some(call) = self.rx.next().await {
250 match call {
251 PickerCall::RegisterSources { reply } => {
252 let _ = reply.send(self.call_register_sources().await);
253 }
254 PickerCall::Init {
255 source,
256 ctx,
257 args,
258 reply,
259 } => {
260 let _ = reply.send(self.call_init(&source, &ctx, &args).await);
261 }
262 PickerCall::Accept {
263 source,
264 ctx,
265 routing,
266 reply,
267 } => {
268 let _ = reply.send(self.call_accept(&source, &ctx, &routing).await);
269 }
270 }
271 }
272 }
273
274 /// OR.5b: drive `register-picker-sources`, then drain what the guest
275 /// declared through the imported `register-picker-source`.
276 ///
277 /// The drain reads `PluginState` AFTER the export returns, which is the
278 /// `register-grammar` shape — a guest registers by calling, so the specs do
279 /// not exist until its body has run.
280 async fn call_register_sources(&mut self) -> CallResult<Vec<NativePickerSourceSpec>> {
281 if self.quarantine.is_tripped() {
282 return Err(PluginHostError::Quarantined {
283 func: "register-picker-sources",
284 });
285 }
286 arm_store(&mut self.store, self.budget)?;
287 let __trace_start = std::time::Instant::now();
288 let result = self
289 .bindings
290 .call_register_picker_sources(&mut self.store)
291 .await;
292 crate::trip_and_map_traced(
293 self.tracer.as_ref(),
294 self.id.0,
295 crate::PluginSeam::PickerSource,
296 &mut self.quarantine,
297 "register-picker-sources",
298 __trace_start,
299 result,
300 )?;
301 Ok(self.store.data_mut().picker_contributions.take())
302 }
303
304 async fn call_init(
305 &mut self,
306 source: &str,
307 ctx: &PickerContext,
308 args: &[String],
309 ) -> CallResult<Result<Vec<CandidatePair>, String>> {
310 if self.quarantine.is_tripped() {
311 return Err(PluginHostError::Quarantined { func: "init" });
312 }
313 arm_store(&mut self.store, self.budget)?;
314 let __trace_start = std::time::Instant::now();
315 let result = self
316 .bindings
317 .lattice_plugin_host_picker_source()
318 .call_init(&mut self.store, source, ctx, args)
319 .await;
320 crate::trip_and_map_traced(
321 self.tracer.as_ref(),
322 self.id.0,
323 crate::PluginSeam::PickerSource,
324 &mut self.quarantine,
325 "init",
326 __trace_start,
327 result,
328 )
329 }
330
331 async fn call_accept(
332 &mut self,
333 source: &str,
334 ctx: &PickerContext,
335 routing: &RoutingPayload,
336 ) -> CallResult<Result<PickerAcceptOutcome, String>> {
337 if self.quarantine.is_tripped() {
338 return Err(PluginHostError::Quarantined { func: "accept" });
339 }
340 arm_store(&mut self.store, self.budget)?;
341 let __trace_start = std::time::Instant::now();
342 let result = self
343 .bindings
344 .lattice_plugin_host_picker_source()
345 .call_accept(&mut self.store, source, ctx, routing)
346 .await;
347 crate::trip_and_map_traced(
348 self.tracer.as_ref(),
349 self.id.0,
350 crate::PluginSeam::PickerSource,
351 &mut self.quarantine,
352 "accept",
353 __trace_start,
354 result,
355 )
356 }
357}
358
359impl PluginHost {
360 /// Instantiate a `picker-source-plugin` component under its capability grant
361 /// and return the bridge: a `Send + Sync` [`PickerClient`] plus the
362 /// [`PickerActor`] the caller drives (spawn [`PickerActor::run`] on a
363 /// multi-thread runtime). Grant computation, the private data dir, and the
364 /// scoped WASI view are identical to
365 /// [`instantiate_plugin`](Self::instantiate_plugin) (via `build_plugin_wasi`)
366 /// — a picker plugin is sandboxed exactly like a lifecycle plugin.
367 ///
368 /// The actor is *not* spawned here (the lib owns no runtime). Until the
369 /// caller drives it, calls on the client simply queue on the channel.
370 ///
371 /// Denied capabilities (a tier-withheld request) are logged; surfacing them
372 /// to the user rides the registration path (PH7.4c.2), which is the only
373 /// consumer that needs them.
374 pub async fn spawn_picker_source(
375 &self,
376 component: &Component,
377 manifest: &PluginManifest,
378 tier: TrustTier,
379 budget: PluginBudget,
380 bus: &Arc<EventBus>,
381 config: Option<&Arc<lattice_config::ConfigRegistry>>,
382 ) -> Result<(PickerClient, PickerActor), PluginHostError> {
383 let (wasi, outcome, _data_dir) = self.build_plugin_wasi(manifest, tier);
384 for denied in &outcome.denied {
385 tracing::warn!(
386 plugin = %manifest.id,
387 capability = ?denied,
388 "picker plugin loaded with a withheld capability (reduced function)"
389 );
390 }
391 let mut store = self.new_store(wasi, outcome.grant, budget, Some(&manifest.id))?;
392 let bindings = PickerSourcePlugin::instantiate_async(&mut store, component, &self.linker)
393 .await
394 .map_err(|e| PluginHostError::Instantiate(e.into()))?;
395 let id = self.alloc_id();
396 // PO.5: route this plugin's `logging` calls into the tracer (Layer 2).
397 store.data_mut().log_ctx = self.log_ctx_for(id);
398 // OR.6: the config registry, so a source can read the options that
399 // decide what it offers.
400 //
401 // Its absence was not a gap in the abstract — org-roam's find-node reads
402 // `org.roam-directory` to decide whether it is configured at all, and
403 // without a registry on THIS store `get-option` answered `none` and the
404 // picker reported "roam is not configured" for a corpus it had just
405 // indexed. The seam was wired end to end and answered nothing, which is
406 // the failure `spawn_event_plugin` / `spawn_context_plugin` /
407 // `spawn_transient_plugin` each already carry this line to prevent.
408 if let Some(registry) = config {
409 store.data_mut().config_registry = Some(Arc::clone(registry));
410 }
411 let (tx, rx) = mpsc::unbounded();
412 let client = PickerClient { tx, id };
413 let actor = PickerActor {
414 store,
415 bindings,
416 budget,
417 rx,
418 id,
419 quarantine: crate::Quarantine::new(id, Arc::clone(bus)),
420 tracer: None,
421 };
422 Ok((client, actor))
423 }
424}