frust/image_async.rs
1//! Async image decode: [`decode_image_async`] is Frust's
2//! off-thread counterpart to the existing synchronous
3//! [`ImageSource::decode`](frust_widgets::ImageSource::decode) — a thin
4//! `frust_reactive::spawn_blocking` wrapper meant to compose with
5//! [`use_task`](frust_reactive::use_task):
6//!
7//! ```ignore
8//! let image = use_task(|| async { frust::decode_image_async(bytes).await });
9//! ```
10//!
11//! # Why this lives in the facade, not `frust-widgets`
12//!
13//! `frust-widgets` is reactive-free by charter (`docs/ARCHITECTURE.md`'s Layer
14//! Dependencies: `frust-widgets = core + scene + text + theme`, no
15//! `frust-reactive`/tokio edge) — adding an async decode entry point there
16//! would pull `frust-reactive` (and transitively tokio) into a crate that
17//! must stay usable with no reactive runtime at all (bare-core tests, a
18//! pre-reactive app). The synchronous `ImageSource::decode` stays exactly
19//! where it is, unchanged, for that reason and for embedded-asset callers
20//! that have no need to leave the calling thread. `decode_image_async` is a
21//! plain function over the facade's own `frust-reactive`/`frust-widgets`
22//! dependencies — no new crate, no widget-side change.
23//!
24//! # The Arc move
25//!
26//! [`ImageSource`] is already an `Arc`-backed handle
27//! (`frust-widgets::image`'s module docs): decoding happens exactly once,
28//! inside the `spawn_blocking` closure, and the resulting `ImageSource`
29//! *moves* back across the thread boundary — the returned value on the UI
30//! thread shares the same underlying `Arc<peniko::ImageData>` allocation the
31//! background thread decoded into, observable via
32//! [`ImageSource::same`](frust_widgets::ImageSource::same) (an `Arc` pointer
33//! check, never a pixel comparison). No second decode, no pixel copy.
34
35use frust_reactive::spawn_blocking;
36use frust_widgets::{ImageError, ImageSource};
37
38/// The error [`decode_image_async`] can fail with.
39///
40/// Two distinct failure modes, kept apart rather than collapsed into one
41/// message: the synchronous decode itself can fail
42/// ([`Decode`](Self::Decode), wrapping the existing
43/// [`ImageError`](frust_widgets::ImageError) unchanged), or the background
44/// `spawn_blocking` task can fail to *run* to completion at all
45/// ([`TaskFailed`](Self::TaskFailed) — a panic inside the decode closure, or
46/// the task being dropped/aborted before it finished; mirrors
47/// `tokio::task::JoinError`'s own two cases). The join-failure text is kept
48/// as a plain `String` here rather than naming `tokio::task::JoinError`
49/// directly: `frust` has no *direct* dependency on `tokio` (only a
50/// transitive one through `frust-reactive`), and this is the one place that
51/// would otherwise need one just to name the type.
52#[derive(Debug)]
53pub enum ImageDecodeError {
54 /// `ImageSource::decode` ran to completion and reported a decode
55 /// failure (bad bytes, unsupported format, ...).
56 Decode(ImageError),
57 /// The background decode task itself panicked or was cancelled before it
58 /// could report a result.
59 TaskFailed(String),
60}
61
62impl std::fmt::Display for ImageDecodeError {
63 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
64 match self {
65 ImageDecodeError::Decode(e) => write!(f, "{e}"),
66 ImageDecodeError::TaskFailed(msg) => write!(f, "image decode task failed: {msg}"),
67 }
68 }
69}
70
71impl std::error::Error for ImageDecodeError {
72 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
73 match self {
74 ImageDecodeError::Decode(e) => Some(e),
75 ImageDecodeError::TaskFailed(_) => None,
76 }
77 }
78}
79
80/// Decodes PNG/JPEG `bytes` off the UI thread via
81/// [`frust_reactive::spawn_blocking`], returning the same
82/// [`ImageSource`](frust_widgets::ImageSource) [`ImageSource::decode`] would
83/// build synchronously.
84///
85/// `bytes` moves into the background closure (no copy beyond the initial
86/// `Vec` transfer); the decoded [`ImageSource`] moves back — see this
87/// module's docs for the Arc-move contract.
88///
89/// Compose with [`use_task`](frust_reactive::use_task) for the full
90/// load/error/ready idiom:
91///
92/// ```ignore
93/// let image = use_task(|| async { frust::decode_image_async(bytes.clone()).await });
94/// ```
95pub async fn decode_image_async(bytes: Vec<u8>) -> Result<ImageSource, ImageDecodeError> {
96 decode_with_probe(
97 bytes,
98 |_thread_id, _result: Result<&ImageSource, &ImageError>| {},
99 )
100 .await
101}
102
103/// The shared decode entry both [`decode_image_async`] and this module's
104/// tests call: identical to `decode_image_async`, except a `probe` closure
105/// runs *inside* the `spawn_blocking` closure, immediately after decoding and
106/// before the result crosses back to the caller — the hooked/spied path the
107/// tests below use to observe the executing thread id and to clone the
108/// decoded `ImageSource`'s `Arc` before the cross, so it can be compared by
109/// pointer identity against the value that arrives on the calling thread.
110async fn decode_with_probe<P>(bytes: Vec<u8>, probe: P) -> Result<ImageSource, ImageDecodeError>
111where
112 P: FnOnce(std::thread::ThreadId, Result<&ImageSource, &ImageError>) + Send + 'static,
113{
114 let outcome = spawn_blocking(move || {
115 let result = ImageSource::decode(&bytes);
116 probe(std::thread::current().id(), result.as_ref());
117 result
118 })
119 .await;
120
121 match outcome {
122 Ok(inner) => inner.map_err(ImageDecodeError::Decode),
123 Err(join_err) => Err(ImageDecodeError::TaskFailed(join_err.to_string())),
124 }
125}
126
127#[cfg(test)]
128mod tests {
129 use super::*;
130 use frust_reactive::{Owner, ReactiveRuntime, use_task};
131 use reactive_graph::traits::GetUntracked;
132 use std::sync::{Arc, Mutex};
133 use std::time::{Duration, Instant};
134
135 /// Pump the UI-thread local queue until `cond` holds or the deadline
136 /// passes; returns whether `cond` became true. Mirrors
137 /// `frust-reactive::task`'s own test helper — the background
138 /// `spawn_blocking` work runs on the reactive runtime's own worker
139 /// threads regardless; this only drains the UI-side coordinator that
140 /// awaits it.
141 fn pump_until(rt: &ReactiveRuntime, timeout: Duration, mut cond: impl FnMut() -> bool) -> bool {
142 let start = Instant::now();
143 loop {
144 rt.pump_local();
145 if cond() {
146 return true;
147 }
148 if start.elapsed() >= timeout {
149 return false;
150 }
151 std::thread::sleep(Duration::from_millis(1));
152 }
153 }
154
155 /// A minimal, hand-encoded 2x2 RGBA PNG, built via the `image` crate's own
156 /// encoder so this test has no external fixture file to keep in sync
157 /// (mirrors `frust-widgets::image`'s own test fixture).
158 fn two_by_two_png() -> Vec<u8> {
159 use image::{ImageEncoder, codecs::png::PngEncoder};
160 let pixels: [u8; 2 * 2 * 4] = [
161 255, 0, 0, 255, // red
162 0, 255, 0, 255, // green
163 0, 0, 255, 255, // blue
164 255, 255, 0, 255, // yellow
165 ];
166 let mut out = Vec::new();
167 PngEncoder::new(&mut out)
168 .write_image(&pixels, 2, 2, image::ExtendedColorType::Rgba8)
169 .expect("encoding a tiny in-memory PNG must not fail");
170 out
171 }
172
173 /// Acceptance criterion 1 + 2: `decode_image_async`'s underlying decode
174 /// runs off the calling (UI) thread, the decoded `ImageSource` crosses
175 /// back as an `Arc` move (same allocation observable via
176 /// `ImageSource::same` on a clone taken before/after the cross), and the
177 /// whole thing composes with `use_task`.
178 #[test]
179 fn decode_runs_off_ui_thread_and_result_is_an_arc_move_via_use_task() {
180 let rt = ReactiveRuntime::init(Arc::new(|| {}));
181 let ui_thread = std::thread::current().id();
182
183 let observed_thread: Arc<Mutex<Option<std::thread::ThreadId>>> = Arc::new(Mutex::new(None));
184 // Cloned *inside* the spawn_blocking closure, before the result
185 // crosses back to the UI thread — the "before" half of the
186 // before/after Arc-identity check.
187 let before_cross: Arc<Mutex<Option<ImageSource>>> = Arc::new(Mutex::new(None));
188
189 let bytes = two_by_two_png();
190 let owner = Owner::new();
191 let task = {
192 let observed_thread = observed_thread.clone();
193 let before_cross = before_cross.clone();
194 owner.with(|| {
195 use_task(move || {
196 let observed_thread = observed_thread.clone();
197 let before_cross = before_cross.clone();
198 let bytes = bytes.clone();
199 async move {
200 decode_with_probe(bytes, move |thread_id, result| {
201 *observed_thread.lock().expect("observed_thread poisoned") =
202 Some(thread_id);
203 if let Ok(source) = result {
204 *before_cross.lock().expect("before_cross poisoned") =
205 Some(source.clone());
206 }
207 })
208 .await
209 }
210 })
211 })
212 };
213
214 assert!(
215 pump_until(rt, Duration::from_secs(5), || task
216 .signal()
217 .get_untracked()
218 .is_ready()),
219 "decode_image_async must resolve to Ready via use_task"
220 );
221
222 let decoded_thread = observed_thread
223 .lock()
224 .expect("observed_thread poisoned")
225 .expect("probe must have recorded a thread id");
226 assert_ne!(
227 decoded_thread, ui_thread,
228 "the decode must run off the calling (UI) thread"
229 );
230
231 let after_cross = task
232 .signal()
233 .get_untracked()
234 .ready()
235 .cloned()
236 .expect("Ready must carry the decoded ImageSource");
237 let before_cross = before_cross
238 .lock()
239 .expect("before_cross poisoned")
240 .clone()
241 .expect("probe must have captured a pre-cross clone");
242 assert!(
243 before_cross.same(&after_cross),
244 "the ImageSource observed inside spawn_blocking (pre-cross) must share the \
245 same Arc allocation as the value observed on the UI thread (post-cross) — \
246 the result crosses as an Arc move, never a re-decode or pixel copy"
247 );
248
249 owner.cleanup();
250 }
251
252 /// The ordinary happy path, exactly at the call site the module docs
253 /// promise (`use_task(|| async { frust::decode_image_async(bytes).await
254 /// })`), with no probe involved.
255 #[test]
256 fn decode_image_async_composes_with_use_task() {
257 let rt = ReactiveRuntime::init(Arc::new(|| {}));
258 let bytes = two_by_two_png();
259 let owner = Owner::new();
260 let task = owner.with(|| {
261 use_task(move || {
262 let bytes = bytes.clone();
263 async move { decode_image_async(bytes).await }
264 })
265 });
266
267 assert!(
268 pump_until(rt, Duration::from_secs(5), || task
269 .signal()
270 .get_untracked()
271 .is_ready()),
272 "a valid PNG must decode to Ready"
273 );
274 let source = task.signal().get_untracked().ready().cloned().unwrap();
275 assert_eq!(source.natural_size(), kurbo::Size::new(2.0, 2.0));
276
277 owner.cleanup();
278 }
279
280 /// A decode failure surfaces through `use_task`'s `Error` state, carrying
281 /// `ImageDecodeError` — proving the error type composes with `use_task`'s
282 /// `E: std::error::Error + Send + Sync + 'static` bound.
283 #[test]
284 fn decode_failure_surfaces_as_use_task_error() {
285 let rt = ReactiveRuntime::init(Arc::new(|| {}));
286 let owner = Owner::new();
287 let task = owner
288 .with(|| use_task(|| async { decode_image_async(b"not an image".to_vec()).await }));
289
290 assert!(
291 pump_until(rt, Duration::from_secs(5), || task
292 .signal()
293 .get_untracked()
294 .is_error()),
295 "invalid bytes must resolve to Error"
296 );
297 assert!(matches!(
298 task.signal().get_untracked().error().map(|e| e.to_string()),
299 Some(msg) if msg.contains("failed to decode image")
300 ));
301
302 owner.cleanup();
303 }
304}