Skip to main content

lattice_plugin_host/
media_task.rs

1//! IM.6b — the per-plugin actor bridge for inline-media providers.
2//!
3//! The media analogue of `decoration_task.rs`, and deliberately its near-twin:
4//! a dedicated async task owns the plugin's `Store<PluginState>` for life (the
5//! Store is `!Sync`), a [`MediaCall`] crosses an mpsc channel with a `oneshot`
6//! reply, and the `Send + Sync` [`MediaClient`] serialises calls onto the
7//! single-consumer loop.
8//!
9//! A **producer**, host-called OFF the render path: the host calls `produce` on
10//! a trigger (buffer opened, edited), caches the result, and the renderer reads
11//! the cache. WASM on the tick would be a paramount-#1 violation.
12//!
13//! Generalising the picker / completion / decoration / media actors over their
14//! bindings type is still deferred — this is the fourth near-copy, so the
15//! rule-of-three note in `completion_task` has now been earned and the
16//! generalisation is worth doing the next time one of them changes shape.
17
18use std::sync::Arc;
19
20use futures::StreamExt;
21use futures::channel::{mpsc, oneshot};
22use lattice_runtime::EventBus;
23use wasmtime::Store;
24
25use crate::media_host::bindings::MediaPlugin;
26use crate::{
27    Component, PluginBudget, PluginHost, PluginHostError, PluginId, PluginManifest, PluginState,
28    TrustTier, arm_store,
29};
30
31// The `with:`-mapped mirrors — the SAME Rust types the boundary round-trips.
32pub use crate::lattice::plugin_host::types::{DecorationContext, MediaBlock, MediaFit};
33
34type CallResult<T> = Result<T, PluginHostError>;
35
36enum MediaCall {
37    /// `media.media-blocks(ctx)` — produce the buffer's inline media blocks.
38    Produce {
39        ctx: Box<DecorationContext>,
40        text: String,
41        reply: oneshot::Sender<CallResult<Result<Vec<MediaBlock>, String>>>,
42    },
43}
44
45/// The `Send + Sync` handle a caller holds. Cloning is cheap; every clone talks
46/// to the same actor / `Store`, so calls serialise on the single-consumer loop
47/// the `!Sync` `Store` requires.
48#[derive(Clone, Debug)]
49pub struct MediaClient {
50    tx: mpsc::UnboundedSender<MediaCall>,
51    id: PluginId,
52}
53
54impl MediaClient {
55    pub fn id(&self) -> PluginId {
56        self.id
57    }
58
59    /// Call the guest's `media-blocks(ctx)`.
60    ///
61    /// The outer result is the host surface (trap / gone / quarantined); the
62    /// inner `Result<_, String>` is the guest's own WIT `result`. An `Err`
63    /// string means the provider produced nothing this trigger — it is logged
64    /// and the buffer KEEPS its prior blocks, so a transient failure mid-edit
65    /// does not make every image in the document blink out.
66    pub async fn produce(
67        &self,
68        ctx: DecorationContext,
69        text: String,
70    ) -> CallResult<Result<Vec<MediaBlock>, String>> {
71        let (reply, rx) = oneshot::channel();
72        self.tx
73            .unbounded_send(MediaCall::Produce {
74                ctx: Box::new(ctx),
75                text,
76                reply,
77            })
78            .map_err(|_| PluginHostError::PluginGone {
79                func: "media-blocks",
80            })?;
81        rx.await.map_err(|_| PluginHostError::PluginGone {
82            func: "media-blocks",
83        })?
84    }
85}
86
87/// The per-plugin actor: owns the `Store` + media bindings for the plugin's
88/// life and serves calls off the channel until every [`MediaClient`] drops.
89pub struct MediaActor {
90    store: Store<PluginState>,
91    bindings: MediaPlugin,
92    budget: PluginBudget,
93    rx: mpsc::UnboundedReceiver<MediaCall>,
94    id: PluginId,
95    /// Crash-quarantine: the first `media-blocks` trap trips this, fires one
96    /// `PluginCrashed`, and every later call returns `Quarantined`.
97    quarantine: crate::Quarantine,
98    tracer: Option<crate::trace::PluginTracerHandle>,
99}
100
101impl MediaActor {
102    pub fn id(&self) -> PluginId {
103        self.id
104    }
105
106    pub fn with_tracer(mut self, tracer: Option<crate::trace::PluginTracerHandle>) -> Self {
107        self.tracer = tracer;
108        self
109    }
110
111    pub async fn run(mut self) {
112        while let Some(call) = self.rx.next().await {
113            match call {
114                MediaCall::Produce { ctx, text, reply } => {
115                    let _ = reply.send(self.call_produce(&ctx, &text).await);
116                }
117            }
118        }
119    }
120
121    async fn call_produce(
122        &mut self,
123        ctx: &DecorationContext,
124        text: &str,
125    ) -> CallResult<Result<Vec<MediaBlock>, String>> {
126        if self.quarantine.is_tripped() {
127            return Err(PluginHostError::Quarantined {
128                func: "media-blocks",
129            });
130        }
131        arm_store(&mut self.store, self.budget)?;
132        let __trace_start = std::time::Instant::now();
133        let result = self
134            .bindings
135            .lattice_plugin_host_media()
136            .call_media_blocks(&mut self.store, ctx, text)
137            .await;
138        crate::trip_and_map_traced(
139            self.tracer.as_ref(),
140            self.id.0,
141            crate::PluginSeam::Media,
142            &mut self.quarantine,
143            "media-blocks",
144            __trace_start,
145            result,
146        )
147    }
148}
149
150impl PluginHost {
151    /// Instantiate a `media-plugin` component under its capability grant and
152    /// return the bridge. Grant / data-dir / WASI are identical to every other
153    /// seam (shared `build_plugin_wasi` + `new_store`), and the actor is NOT
154    /// spawned here — the lib owns no runtime.
155    pub async fn spawn_media_source(
156        &self,
157        component: &Component,
158        manifest: &PluginManifest,
159        tier: TrustTier,
160        budget: PluginBudget,
161        bus: &Arc<EventBus>,
162    ) -> Result<(MediaClient, MediaActor), PluginHostError> {
163        let (wasi, outcome, _data_dir) = self.build_plugin_wasi(manifest, tier);
164        for denied in &outcome.denied {
165            tracing::warn!(
166                plugin = %manifest.id,
167                capability = ?denied,
168                "media plugin loaded with a withheld capability (reduced function)"
169            );
170        }
171        let mut store = self.new_store(wasi, outcome.grant, budget, Some(&manifest.id))?;
172        let bindings = MediaPlugin::instantiate_async(&mut store, component, &self.linker)
173            .await
174            .map_err(|e| PluginHostError::Instantiate(e.into()))?;
175        let id = self.alloc_id();
176        store.data_mut().log_ctx = self.log_ctx_for(id);
177        let (tx, rx) = mpsc::unbounded();
178        let client = MediaClient { tx, id };
179        let actor = MediaActor {
180            store,
181            bindings,
182            budget,
183            rx,
184            id,
185            quarantine: crate::Quarantine::new(id, Arc::clone(bus)),
186            tracer: None,
187        };
188        Ok((client, actor))
189    }
190}