1use 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
31pub use crate::lattice::plugin_host::types::{DecorationContext, MediaBlock, MediaFit};
33
34type CallResult<T> = Result<T, PluginHostError>;
35
36enum MediaCall {
37 Produce {
39 ctx: Box<DecorationContext>,
40 text: String,
41 reply: oneshot::Sender<CallResult<Result<Vec<MediaBlock>, String>>>,
42 },
43}
44
45#[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 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
87pub struct MediaActor {
90 store: Store<PluginState>,
91 bindings: MediaPlugin,
92 budget: PluginBudget,
93 rx: mpsc::UnboundedReceiver<MediaCall>,
94 id: PluginId,
95 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 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}