Skip to main content

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}