Skip to main content

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}