Skip to main content

lattice_plugin_host/
scan_task.rs

1//! OM.A1 — the per-plugin actor bridge for scanned-excerpt-source providers.
2//!
3//! The agenda analogue of `media_task.rs`, and deliberately its near-twin: a
4//! dedicated async task owns the plugin's `Store<PluginState>` for life (the
5//! Store is `!Sync`), an [`ScanCall`] crosses an mpsc channel with a
6//! `oneshot` reply, and the `Send + Sync` [`ScanClient`] serialises calls
7//! onto the single-consumer loop.
8//!
9//! **The serialisation is load-bearing here, not incidental.** `begin` drops
10//! per-scan state and every following `scan` reads it, so the two must not
11//! interleave with a second scan's calls. One actor, one queue, in order —
12//! which is also why the scan walks files sequentially rather than fanning
13//! out across the pool.
14//!
15//! This is the fifth near-copy of the picker / completion / decoration / media
16//! actor. The rule-of-three note in `completion_task` is now well past earned;
17//! generalising over the bindings type is worth doing the next time one of
18//! them changes shape, and this slice deliberately did not take that on
19//! mid-seam.
20
21use 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
34/// HB.5: the row an [`Entry`] hangs below itself.
35pub use crate::scan_host::bindings::lattice::plugin_host::scanned_excerpt_source::Annotation;
36/// The WIT `entry`, re-exported so the adapter and the loader name one type.
37pub use crate::scan_host::bindings::lattice::plugin_host::scanned_excerpt_source::Entry;
38/// OA.14b: a file's clocked time, reported beside its rows.
39pub use crate::scan_host::bindings::lattice::plugin_host::scanned_excerpt_source::{
40    ClockSpan, ScanResult,
41};
42/// OA.5: the per-row style spans an [`Entry`] carries.
43pub use crate::scan_host::bindings::lattice::plugin_host::types::DisplaySpan;
44
45/// OT.3: parse one scanned file, from text the host has already read.
46///
47/// `None` means "hand the guest the text instead" — an extension that resolves
48/// to no registered language, a grammar that will not load, or a parse that
49/// yields no tree. None of those is an error: `scanned-excerpt-source.wit` keeps a
50/// source independent of the `language` seam, so a filetype with no grammar
51/// must still be scannable.
52///
53/// `debug!`, never `info!` — this runs once per file of a project-wide walk,
54/// and a tree of unparseable files is the ordinary case rather than a problem.
55fn 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()` — the file extensions this source wants offered.
81    Extensions {
82        reply: oneshot::Sender<CallResult<Vec<String>>>,
83    },
84    /// `view-mode()` — the minor this source wants on the agenda view.
85    ViewMode {
86        reply: oneshot::Sender<CallResult<Option<String>>>,
87    },
88    /// AF.1: `roots()` — the paths this source wants scanned.
89    Roots {
90        reply: oneshot::Sender<CallResult<Vec<String>>>,
91    },
92    /// `begin(args)` — drop per-scan state, and return the generation key
93    /// (OT.3b). `args` (OA.11a) are the view's own scan arguments, which the
94    /// host forwards without reading.
95    Begin {
96        args: Vec<String>,
97        reply: oneshot::Sender<CallResult<u64>>,
98    },
99    /// OA.22: `describe(args)` — what this view is, for its headerline.
100    Describe {
101        args: Vec<String>,
102        reply: oneshot::Sender<CallResult<String>>,
103    },
104    /// `scan(path, text)` — one file's agenda rows.
105    Scan {
106        path: String,
107        text: String,
108        reply: oneshot::Sender<CallResult<Result<ScanResult, String>>>,
109    },
110}
111
112/// The `Send + Sync` handle a caller holds. Cloning is cheap; every clone
113/// talks to the same actor / `Store`, so calls serialise on the
114/// single-consumer loop the `!Sync` `Store` requires — which is exactly what
115/// `begin`-then-`scan` needs.
116#[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    /// Call the guest's `extensions()`. Once, at load.
128    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    /// Call the guest's `view-mode()`. Once, at load.
138    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    /// AF.1: call the guest's `roots()`. Per scan, not once at load — the
148    /// answer comes from user configuration and must follow a `:set`.
149    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    /// Call the guest's `begin()` — the start of a scan. Returns the guest's
159    /// generation key (OT.3b): an opaque `u64` that changes when anything
160    /// scan-wide would change its rows, so the host can invalidate cached
161    /// results without knowing what those things are.
162    ///
163    /// OA.11a: `args` are the view's own scan arguments, forwarded verbatim.
164    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    /// OA.22: ask the guest what this view IS, for the headerline.
174    ///
175    /// Called once per scan, after [`Self::begin`], so the guest sees the args
176    /// `begin` stashed. The host cannot answer this itself: it deliberately
177    /// does not read `args`, so a filtered agenda would otherwise look exactly
178    /// like an unfiltered one.
179    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    /// Call the guest's `scan(path, text)`.
189    ///
190    /// The outer result is the host surface (trap / gone / quarantined); the
191    /// inner `Result<_, String>` is the guest's own WIT `result`. Either way
192    /// the caller skips THIS FILE and continues — one bad file must not fail
193    /// the agenda.
194    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
204/// The per-plugin actor: owns the `Store` + agenda bindings for the plugin's
205/// life and serves calls off the channel until every [`ScanClient`] drops.
206pub struct ScanActor {
207    store: Store<PluginState>,
208    bindings: ScannedExcerptSourcePlugin,
209    budget: PluginBudget,
210    rx: mpsc::UnboundedReceiver<ScanCall>,
211    id: PluginId,
212    /// Crash-quarantine: the first trap trips this, fires one `PluginCrashed`,
213    /// and every later call returns `Quarantined`. A trap mid-scan therefore
214    /// leaves the agenda showing what it collected rather than emptying it —
215    /// partial-and-honest beats empty-and-silent (`org-mode.md` §8).
216    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    /// **Fuel is re-armed per call**, and a scan calls this once per file.
346    /// Arming once at instantiate — correct for a declare-once seam — would
347    /// make the agenda work for the first stretch of a large project and then
348    /// silently contribute nothing for the rest, which is precisely the cliff
349    /// `WasmErrorParser::rearm` documents. `arm_store` above is that re-arm;
350    /// it is not boilerplate.
351    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        // OT.3: parse here, from the text the host already read, and lend the
362        // guest a borrow ALONGSIDE that text. Not instead of it — the copy is
363        // 217 ns (`benches/agenda_scan_input.rs`) while the parse buying the
364        // tree is 1-2 ms, and the seam exposes no node TEXT, so a guest without
365        // the string would need a crossing per headline to read a TODO keyword.
366        // Structure from the tree, characters from the text.
367        //
368        // NOT `tree-sitter.parse-file`: that is `fs:read` gated, and this is the
369        // one seam whose guest holds no capability at all (`scanned-excerpt-source.wit`
370        // — "no preopens, no `walk`"). The host reads; the guest is handed the
371        // result.
372        //
373        // `None` when the extension resolves to no language or the parse yields
374        // no tree; the guest then scans text alone, as it always did — a source
375        // is independent of the `language` seam, so a filetype with no grammar
376        // must still scan.
377        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        // Reclaim the lent entry — the host owns it throughout (the
398        // `apply-action` pattern), and a scan leaking one per file would grow
399        // the resource table for the length of a project walk.
400        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    /// Instantiate an `scanned-excerpt-source-plugin` component under its capability
417    /// grant and return the bridge. Grant / data-dir / WASI are identical to
418    /// every other seam (shared `build_plugin_wasi` + `new_store`), and the
419    /// actor is NOT spawned here — the lib owns no runtime.
420    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        // AF.3: the editor's option registry, so `config.get-option` answers
428        // inside a `roots` / `begin` / `scan` call.
429        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        // AF.3: without this, `config.get-option` answers `none` in every
447        // agenda call and the guest silently falls back to its compiled
448        // defaults — org's `agenda-files` would read as unset however the user
449        // set it, and its TODO keywords would ignore `org.todo-keywords`
450        // entirely.
451        //
452        // `context`, `event`, `transient` and `grammar` all wire this; the
453        // agenda store did not, and 73842466 fixed exactly this omission for
454        // the events store. It survived here because every agenda test drove
455        // `extensions` / `begin` / `scan`, none of which read an option until
456        // `roots` did — a seam covered end to end with its config path never
457        // called once.
458        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}