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}