Skip to main content

lattice_plugin_host/
completion_task.rs

1//! PH7.6 — the per-plugin actor bridge for completion sources.
2//!
3//! The completion analogue of `picker_task.rs`: a dedicated async task owns the
4//! plugin's `Store<PluginState>` for life (the Store is `!Sync`), a
5//! `CompletionCall` crosses an mpsc channel with a `oneshot` reply, and the
6//! `Send + Sync` [`CompletionClient`] serializes calls onto the single-consumer
7//! loop. `PluginHost::spawn_completion_source` instantiates the
8//! `completion-source-plugin` world under the plugin's grant and returns
9//! `(CompletionClient, CompletionActor)`; the caller drives
10//! [`CompletionActor::run`] on its multi-thread runtime (the lib owns no
11//! runtime).
12//!
13//! The picker and completion actors are near-identical; a third guest-backed
14//! `Arc<dyn>` adapter (grammar, PH7.7) is the rule-of-three trigger to generalise
15//! the loop over the bindings type. Until then this duplicates the ~50-line
16//! shape and reuses the shared `arm_store` / `new_store` / `build_plugin_wasi`
17//! primitives.
18
19use std::sync::Arc;
20
21use futures::StreamExt;
22use futures::channel::{mpsc, oneshot};
23use lattice_runtime::EventBus;
24use wasmtime::Store;
25
26use crate::completion_host::bindings::CompletionSourcePlugin;
27use crate::{
28    Component, PluginBudget, PluginHost, PluginHostError, PluginId, PluginManifest, PluginState,
29    TrustTier, arm_store,
30};
31
32// The completion WIT records the bridge's public API traffics in — the
33// `with:`-mapped `types` mirrors (`completion_host.rs`), i.e. the SAME Rust types
34// `WitBoundary` round-trips; the native↔WIT conversion is the adapter's job
35// (`picker_source.rs` precedent). Re-exported `pub` (they appear in the
36// `CompletionClient` method signatures).
37pub use crate::lattice::plugin_host::types::{CompletionSourceSpec, GenerateContext, RawCandidate};
38
39/// See `picker_task::CallResult`.
40type CallResult<T> = Result<T, PluginHostError>;
41
42/// A request sent from a [`CompletionClient`] to its [`CompletionActor`].
43enum CompletionCall {
44    /// `completion-source.spec()` — the source's `(id, doc)` identity. No WIT
45    /// `result`; the reply is the spec or a host-side trap.
46    Spec {
47        reply: oneshot::Sender<CallResult<CompletionSourceSpec>>,
48    },
49    /// `completion-source.generate(ctx)` — produce raw candidates. Replies the
50    /// guest's `result<list<raw-candidate>, string>` (or a host trap).
51    Generate {
52        ctx: Box<GenerateContext>,
53        reply: oneshot::Sender<CallResult<Result<Vec<RawCandidate>, String>>>,
54    },
55}
56
57/// The `Send + Sync` handle the [`WasmCompletionSource`](crate::WasmCompletionSource)
58/// adapter holds. Cloning is cheap (an mpsc `Sender` clone); every clone talks to
59/// the same actor / `Store`, so calls serialize on the single-consumer loop.
60#[derive(Clone)]
61pub struct CompletionClient {
62    tx: mpsc::UnboundedSender<CompletionCall>,
63    id: PluginId,
64}
65
66impl CompletionClient {
67    /// The host-issued identity of the plugin behind this client.
68    pub fn id(&self) -> PluginId {
69        self.id
70    }
71
72    /// Call the guest's `spec()`.
73    pub async fn spec(&self) -> CallResult<CompletionSourceSpec> {
74        let (reply, rx) = oneshot::channel();
75        self.dispatch(CompletionCall::Spec { reply }, rx, "spec")
76            .await
77    }
78
79    /// Call the guest's `generate(ctx)`. The outer result is the host surface;
80    /// the inner `Result<_, String>` is the guest's own WIT `result` (an `Err`
81    /// string is a source that produced no rows — logged, echoed as empty).
82    pub async fn generate(
83        &self,
84        ctx: GenerateContext,
85    ) -> CallResult<Result<Vec<RawCandidate>, String>> {
86        let (reply, rx) = oneshot::channel();
87        self.dispatch(
88            CompletionCall::Generate {
89                ctx: Box::new(ctx),
90                reply,
91            },
92            rx,
93            "generate",
94        )
95        .await
96    }
97
98    async fn dispatch<T>(
99        &self,
100        call: CompletionCall,
101        rx: oneshot::Receiver<CallResult<T>>,
102        func: &'static str,
103    ) -> CallResult<T> {
104        self.tx
105            .unbounded_send(call)
106            .map_err(|_| PluginHostError::PluginGone { func })?;
107        rx.await.map_err(|_| PluginHostError::PluginGone { func })?
108    }
109}
110
111/// The per-plugin actor: owns the `Store` + completion bindings for the plugin's
112/// life and serves calls off the channel until every [`CompletionClient`] drops.
113pub struct CompletionActor {
114    store: Store<PluginState>,
115    bindings: CompletionSourcePlugin,
116    budget: PluginBudget,
117    rx: mpsc::UnboundedReceiver<CompletionCall>,
118    id: PluginId,
119    /// Crash-quarantine (PH7.12): the first export trap trips this, fires one
120    /// `PluginCrashed`, and every later call returns `Quarantined`.
121    quarantine: crate::Quarantine,
122    /// PO.2: the boundary tracer, wired by the loader via with_tracer; None in tests / pre-wire.
123    tracer: Option<crate::trace::PluginTracerHandle>,
124}
125
126impl CompletionActor {
127    /// The host-issued identity of this plugin.
128    pub fn id(&self) -> PluginId {
129        self.id
130    }
131
132    /// PO.2: attach the boundary tracer (the loader calls this before spawning
133    /// run()). Off the hot path — the seam is async.
134    pub fn with_tracer(mut self, tracer: Option<crate::trace::PluginTracerHandle>) -> Self {
135        self.tracer = tracer;
136        self
137    }
138
139    /// Drive the actor to completion — see `picker_task::PickerActor::run`.
140    pub async fn run(mut self) {
141        while let Some(call) = self.rx.next().await {
142            match call {
143                CompletionCall::Spec { reply } => {
144                    let _ = reply.send(self.call_spec().await);
145                }
146                CompletionCall::Generate { ctx, reply } => {
147                    let _ = reply.send(self.call_generate(&ctx).await);
148                }
149            }
150        }
151    }
152
153    async fn call_spec(&mut self) -> CallResult<CompletionSourceSpec> {
154        if self.quarantine.is_tripped() {
155            return Err(PluginHostError::Quarantined { func: "spec" });
156        }
157        arm_store(&mut self.store, self.budget)?;
158        let __trace_start = std::time::Instant::now();
159        let result = self
160            .bindings
161            .lattice_plugin_host_completion_source()
162            .call_spec(&mut self.store)
163            .await;
164        crate::trip_and_map_traced(
165            self.tracer.as_ref(),
166            self.id.0,
167            crate::PluginSeam::CompletionSource,
168            &mut self.quarantine,
169            "spec",
170            __trace_start,
171            result,
172        )
173    }
174
175    async fn call_generate(
176        &mut self,
177        ctx: &GenerateContext,
178    ) -> CallResult<Result<Vec<RawCandidate>, String>> {
179        if self.quarantine.is_tripped() {
180            return Err(PluginHostError::Quarantined { func: "generate" });
181        }
182        arm_store(&mut self.store, self.budget)?;
183        let __trace_start = std::time::Instant::now();
184        let result = self
185            .bindings
186            .lattice_plugin_host_completion_source()
187            .call_generate(&mut self.store, ctx)
188            .await;
189        crate::trip_and_map_traced(
190            self.tracer.as_ref(),
191            self.id.0,
192            crate::PluginSeam::CompletionSource,
193            &mut self.quarantine,
194            "generate",
195            __trace_start,
196            result,
197        )
198    }
199}
200
201impl PluginHost {
202    /// Instantiate a `completion-source-plugin` component under its capability
203    /// grant and return the bridge: a `Send + Sync` [`CompletionClient`] plus the
204    /// [`CompletionActor`] the caller drives. Grant / data-dir / WASI are
205    /// identical to `instantiate_plugin` (shared `build_plugin_wasi` +
206    /// `new_store`), and the actor is *not* spawned here (the lib owns no
207    /// runtime). Mirror of [`spawn_picker_source`](Self::spawn_picker_source).
208    pub async fn spawn_completion_source(
209        &self,
210        component: &Component,
211        manifest: &PluginManifest,
212        tier: TrustTier,
213        budget: PluginBudget,
214        bus: &Arc<EventBus>,
215        config: Option<&Arc<lattice_config::ConfigRegistry>>,
216    ) -> Result<(CompletionClient, CompletionActor), PluginHostError> {
217        let (wasi, outcome, _data_dir) = self.build_plugin_wasi(manifest, tier);
218        for denied in &outcome.denied {
219            tracing::warn!(
220                plugin = %manifest.id,
221                capability = ?denied,
222                "completion plugin loaded with a withheld capability (reduced function)"
223            );
224        }
225        let mut store = self.new_store(wasi, outcome.grant, budget, Some(&manifest.id))?;
226        let bindings =
227            CompletionSourcePlugin::instantiate_async(&mut store, component, &self.linker)
228                .await
229                .map_err(|e| PluginHostError::Instantiate(e.into()))?;
230        let id = self.alloc_id();
231        // PO.5: route this plugin's `logging` calls into the tracer (Layer 2).
232        store.data_mut().log_ctx = self.log_ctx_for(id);
233        // OR.7: the config registry — the sixth seam to need this line, and
234        // the sixth to have shipped without it. A completion source reads
235        // options to decide what it offers (org-roam's node source needs
236        // `org.roam-directory`); without a registry on THIS store
237        // `get-option` answers `none` and the source silently offers
238        // nothing, which is indistinguishable from "no matches".
239        if let Some(registry) = config {
240            store.data_mut().config_registry = Some(Arc::clone(registry));
241        }
242        let (tx, rx) = mpsc::unbounded();
243        let client = CompletionClient { tx, id };
244        let actor = CompletionActor {
245            store,
246            bindings,
247            budget,
248            rx,
249            id,
250            quarantine: crate::Quarantine::new(id, Arc::clone(bus)),
251            tracer: None,
252        };
253        Ok((client, actor))
254    }
255}