1use std::sync::Arc;
22
23use futures::StreamExt;
24use futures::channel::{mpsc, oneshot};
25use lattice_runtime::EventBus;
26use wasmtime::Store;
27
28use crate::scan_host::bindings::ScannedExcerptSourcePlugin;
29use crate::{
30 Component, PluginBudget, PluginHost, PluginHostError, PluginId, PluginManifest, PluginState,
31 TrustTier, arm_store,
32};
33
34pub use crate::scan_host::bindings::lattice::plugin_host::scanned_excerpt_source::Annotation;
36pub use crate::scan_host::bindings::lattice::plugin_host::scanned_excerpt_source::Entry;
38pub use crate::scan_host::bindings::lattice::plugin_host::scanned_excerpt_source::{
40 ClockSpan, ScanResult,
41};
42pub use crate::scan_host::bindings::lattice::plugin_host::types::DisplaySpan;
44
45fn parse_for_scan(path: &str, text: &str) -> Option<Arc<lattice_syntax::SyntaxSnapshot>> {
56 let lang = lattice_syntax::Lang::detect_from_path(Some(std::path::Path::new(path)));
57 let mut syntax = match lattice_syntax::Syntax::for_language(lang) {
58 Ok(Some(syntax)) => syntax,
59 Ok(None) => {
60 tracing::debug!(%path, ?lang, "scan: no grammar; handing the guest text");
61 return None;
62 }
63 Err(error) => {
64 tracing::debug!(%path, ?lang, %error, "scan: grammar load failed; handing the guest text");
65 return None;
66 }
67 };
68 syntax.parse(text);
69 let snapshot = Arc::new(syntax.snapshot_owned());
70 if snapshot.tree().is_none() {
71 tracing::debug!(%path, ?lang, "scan: parsed to no tree; handing the guest text");
72 return None;
73 }
74 Some(snapshot)
75}
76
77type CallResult<T> = Result<T, PluginHostError>;
78
79enum ScanCall {
80 Extensions {
82 reply: oneshot::Sender<CallResult<Vec<String>>>,
83 },
84 ViewMode {
86 reply: oneshot::Sender<CallResult<Option<String>>>,
87 },
88 Roots {
90 reply: oneshot::Sender<CallResult<Vec<String>>>,
91 },
92 Begin {
96 args: Vec<String>,
97 reply: oneshot::Sender<CallResult<u64>>,
98 },
99 Describe {
101 args: Vec<String>,
102 reply: oneshot::Sender<CallResult<String>>,
103 },
104 Scan {
106 path: String,
107 text: String,
108 reply: oneshot::Sender<CallResult<Result<ScanResult, String>>>,
109 },
110}
111
112#[derive(Clone, Debug)]
117pub struct ScanClient {
118 tx: mpsc::UnboundedSender<ScanCall>,
119 id: PluginId,
120}
121
122impl ScanClient {
123 pub fn id(&self) -> PluginId {
124 self.id
125 }
126
127 pub async fn extensions(&self) -> CallResult<Vec<String>> {
129 let (reply, rx) = oneshot::channel();
130 self.tx
131 .unbounded_send(ScanCall::Extensions { reply })
132 .map_err(|_| PluginHostError::PluginGone { func: "extensions" })?;
133 rx.await
134 .map_err(|_| PluginHostError::PluginGone { func: "extensions" })?
135 }
136
137 pub async fn view_mode(&self) -> CallResult<Option<String>> {
139 let (reply, rx) = oneshot::channel();
140 self.tx
141 .unbounded_send(ScanCall::ViewMode { reply })
142 .map_err(|_| PluginHostError::PluginGone { func: "view-mode" })?;
143 rx.await
144 .map_err(|_| PluginHostError::PluginGone { func: "view-mode" })?
145 }
146
147 pub async fn roots(&self) -> CallResult<Vec<String>> {
150 let (reply, rx) = oneshot::channel();
151 self.tx
152 .unbounded_send(ScanCall::Roots { reply })
153 .map_err(|_| PluginHostError::PluginGone { func: "roots" })?;
154 rx.await
155 .map_err(|_| PluginHostError::PluginGone { func: "roots" })?
156 }
157
158 pub async fn begin(&self, args: Vec<String>) -> CallResult<u64> {
165 let (reply, rx) = oneshot::channel();
166 self.tx
167 .unbounded_send(ScanCall::Begin { args, reply })
168 .map_err(|_| PluginHostError::PluginGone { func: "begin" })?;
169 rx.await
170 .map_err(|_| PluginHostError::PluginGone { func: "begin" })?
171 }
172
173 pub async fn describe(&self, args: Vec<String>) -> CallResult<String> {
180 let (reply, rx) = oneshot::channel();
181 self.tx
182 .unbounded_send(ScanCall::Describe { args, reply })
183 .map_err(|_| PluginHostError::PluginGone { func: "describe" })?;
184 rx.await
185 .map_err(|_| PluginHostError::PluginGone { func: "describe" })?
186 }
187
188 pub async fn scan(&self, path: String, text: String) -> CallResult<Result<ScanResult, String>> {
195 let (reply, rx) = oneshot::channel();
196 self.tx
197 .unbounded_send(ScanCall::Scan { path, text, reply })
198 .map_err(|_| PluginHostError::PluginGone { func: "scan" })?;
199 rx.await
200 .map_err(|_| PluginHostError::PluginGone { func: "scan" })?
201 }
202}
203
204pub struct ScanActor {
207 store: Store<PluginState>,
208 bindings: ScannedExcerptSourcePlugin,
209 budget: PluginBudget,
210 rx: mpsc::UnboundedReceiver<ScanCall>,
211 id: PluginId,
212 quarantine: crate::Quarantine,
217 tracer: Option<crate::trace::PluginTracerHandle>,
218}
219
220impl ScanActor {
221 pub fn id(&self) -> PluginId {
222 self.id
223 }
224
225 pub fn with_tracer(mut self, tracer: Option<crate::trace::PluginTracerHandle>) -> Self {
226 self.tracer = tracer;
227 self
228 }
229
230 pub async fn run(mut self) {
231 while let Some(call) = self.rx.next().await {
232 match call {
233 ScanCall::Extensions { reply } => {
234 let _ = reply.send(self.call_extensions().await);
235 }
236 ScanCall::ViewMode { reply } => {
237 let _ = reply.send(self.call_view_mode().await);
238 }
239 ScanCall::Roots { reply } => {
240 let _ = reply.send(self.call_roots().await);
241 }
242 ScanCall::Begin { args, reply } => {
243 let _ = reply.send(self.call_begin(&args).await);
244 }
245 ScanCall::Describe { args, reply } => {
246 let _ = reply.send(self.call_describe(&args).await);
247 }
248 ScanCall::Scan { path, text, reply } => {
249 let _ = reply.send(self.call_scan(&path, &text).await);
250 }
251 }
252 }
253 }
254
255 async fn call_extensions(&mut self) -> CallResult<Vec<String>> {
256 if self.quarantine.is_tripped() {
257 return Err(PluginHostError::Quarantined { func: "extensions" });
258 }
259 arm_store(&mut self.store, self.budget)?;
260 let start = std::time::Instant::now();
261 let result = self.bindings.call_extensions(&mut self.store).await;
262 crate::trip_and_map_traced(
263 self.tracer.as_ref(),
264 self.id.0,
265 crate::PluginSeam::ScannedExcerptSource,
266 &mut self.quarantine,
267 "extensions",
268 start,
269 result,
270 )
271 }
272
273 async fn call_view_mode(&mut self) -> CallResult<Option<String>> {
274 if self.quarantine.is_tripped() {
275 return Err(PluginHostError::Quarantined { func: "view-mode" });
276 }
277 arm_store(&mut self.store, self.budget)?;
278 let start = std::time::Instant::now();
279 let result = self.bindings.call_view_mode(&mut self.store).await;
280 crate::trip_and_map_traced(
281 self.tracer.as_ref(),
282 self.id.0,
283 crate::PluginSeam::ScannedExcerptSource,
284 &mut self.quarantine,
285 "view-mode",
286 start,
287 result,
288 )
289 }
290
291 async fn call_roots(&mut self) -> CallResult<Vec<String>> {
292 if self.quarantine.is_tripped() {
293 return Err(PluginHostError::Quarantined { func: "roots" });
294 }
295 arm_store(&mut self.store, self.budget)?;
296 let start = std::time::Instant::now();
297 let result = self.bindings.call_roots(&mut self.store).await;
298 crate::trip_and_map_traced(
299 self.tracer.as_ref(),
300 self.id.0,
301 crate::PluginSeam::ScannedExcerptSource,
302 &mut self.quarantine,
303 "roots",
304 start,
305 result,
306 )
307 }
308
309 async fn call_describe(&mut self, args: &[String]) -> CallResult<String> {
310 if self.quarantine.is_tripped() {
311 return Err(PluginHostError::Quarantined { func: "describe" });
312 }
313 arm_store(&mut self.store, self.budget)?;
314 let start = std::time::Instant::now();
315 let result = self.bindings.call_describe(&mut self.store, args).await;
316 crate::trip_and_map_traced(
317 self.tracer.as_ref(),
318 self.id.0,
319 crate::PluginSeam::ScannedExcerptSource,
320 &mut self.quarantine,
321 "describe",
322 start,
323 result,
324 )
325 }
326
327 async fn call_begin(&mut self, args: &[String]) -> CallResult<u64> {
328 if self.quarantine.is_tripped() {
329 return Err(PluginHostError::Quarantined { func: "begin" });
330 }
331 arm_store(&mut self.store, self.budget)?;
332 let start = std::time::Instant::now();
333 let result = self.bindings.call_begin(&mut self.store, args).await;
334 crate::trip_and_map_traced(
335 self.tracer.as_ref(),
336 self.id.0,
337 crate::PluginSeam::ScannedExcerptSource,
338 &mut self.quarantine,
339 "begin",
340 start,
341 result,
342 )
343 }
344
345 async fn call_scan(
352 &mut self,
353 path: &str,
354 text: &str,
355 ) -> CallResult<Result<ScanResult, String>> {
356 if self.quarantine.is_tripped() {
357 return Err(PluginHostError::Quarantined { func: "scan" });
358 }
359 arm_store(&mut self.store, self.budget)?;
360 let start = std::time::Instant::now();
361 let snapshot = parse_for_scan(path, text);
378 let owned_tree = match &snapshot {
379 Some(snap) => Some(
380 self.store
381 .data_mut()
382 .table
383 .push(crate::tree_resource::TreeSnapshotResource::new(
384 snap.clone(),
385 ))
386 .map_err(|e| PluginHostError::Linker(e.into()))?,
387 ),
388 None => None,
389 };
390 let tree_borrow = owned_tree
391 .as_ref()
392 .map(|owned| wasmtime::component::Resource::new_borrow(owned.rep()));
393 let result = self
394 .bindings
395 .call_scan(&mut self.store, path, text, tree_borrow)
396 .await;
397 if let Some(owned) = owned_tree {
401 let _ = self.store.data_mut().table.delete(owned);
402 }
403 crate::trip_and_map_traced(
404 self.tracer.as_ref(),
405 self.id.0,
406 crate::PluginSeam::ScannedExcerptSource,
407 &mut self.quarantine,
408 "scan",
409 start,
410 result,
411 )
412 }
413}
414
415impl PluginHost {
416 pub async fn spawn_scan_source(
421 &self,
422 component: &Component,
423 manifest: &PluginManifest,
424 tier: TrustTier,
425 budget: PluginBudget,
426 bus: &Arc<EventBus>,
427 config: Option<&Arc<lattice_config::ConfigRegistry>>,
430 ) -> Result<(ScanClient, ScanActor), PluginHostError> {
431 let (wasi, outcome, _data_dir) = self.build_plugin_wasi(manifest, tier);
432 for denied in &outcome.denied {
433 tracing::warn!(
434 plugin = %manifest.id,
435 capability = ?denied,
436 "scan source loaded with a withheld capability (reduced function)"
437 );
438 }
439 let mut store = self.new_store(wasi, outcome.grant, budget, Some(&manifest.id))?;
440 let bindings =
441 ScannedExcerptSourcePlugin::instantiate_async(&mut store, component, &self.linker)
442 .await
443 .map_err(|e| PluginHostError::Instantiate(e.into()))?;
444 let id = self.alloc_id();
445 store.data_mut().log_ctx = self.log_ctx_for(id);
446 if let Some(registry) = config {
459 store.data_mut().config_registry = Some(Arc::clone(registry));
460 }
461 let (tx, rx) = mpsc::unbounded();
462 let client = ScanClient { tx, id };
463 let actor = ScanActor {
464 store,
465 bindings,
466 budget,
467 rx,
468 id,
469 quarantine: crate::Quarantine::new(id, Arc::clone(bus)),
470 tracer: None,
471 };
472 Ok((client, actor))
473 }
474}