use frust_reactive::spawn_blocking;
use frust_widgets::{ImageError, ImageSource};
#[derive(Debug)]
pub enum ImageDecodeError {
Decode(ImageError),
TaskFailed(String),
}
impl std::fmt::Display for ImageDecodeError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ImageDecodeError::Decode(e) => write!(f, "{e}"),
ImageDecodeError::TaskFailed(msg) => write!(f, "image decode task failed: {msg}"),
}
}
}
impl std::error::Error for ImageDecodeError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
ImageDecodeError::Decode(e) => Some(e),
ImageDecodeError::TaskFailed(_) => None,
}
}
}
pub async fn decode_image_async(bytes: Vec<u8>) -> Result<ImageSource, ImageDecodeError> {
decode_with_probe(
bytes,
|_thread_id, _result: Result<&ImageSource, &ImageError>| {},
)
.await
}
async fn decode_with_probe<P>(bytes: Vec<u8>, probe: P) -> Result<ImageSource, ImageDecodeError>
where
P: FnOnce(std::thread::ThreadId, Result<&ImageSource, &ImageError>) + Send + 'static,
{
let outcome = spawn_blocking(move || {
let result = ImageSource::decode(&bytes);
probe(std::thread::current().id(), result.as_ref());
result
})
.await;
match outcome {
Ok(inner) => inner.map_err(ImageDecodeError::Decode),
Err(join_err) => Err(ImageDecodeError::TaskFailed(join_err.to_string())),
}
}
#[cfg(test)]
mod tests {
use super::*;
use frust_reactive::{Owner, ReactiveRuntime, use_task};
use reactive_graph::traits::GetUntracked;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
fn pump_until(rt: &ReactiveRuntime, timeout: Duration, mut cond: impl FnMut() -> bool) -> bool {
let start = Instant::now();
loop {
rt.pump_local();
if cond() {
return true;
}
if start.elapsed() >= timeout {
return false;
}
std::thread::sleep(Duration::from_millis(1));
}
}
fn two_by_two_png() -> Vec<u8> {
use image::{ImageEncoder, codecs::png::PngEncoder};
let pixels: [u8; 2 * 2 * 4] = [
255, 0, 0, 255, 0, 255, 0, 255, 0, 0, 255, 255, 255, 255, 0, 255, ];
let mut out = Vec::new();
PngEncoder::new(&mut out)
.write_image(&pixels, 2, 2, image::ExtendedColorType::Rgba8)
.expect("encoding a tiny in-memory PNG must not fail");
out
}
#[test]
fn decode_runs_off_ui_thread_and_result_is_an_arc_move_via_use_task() {
let rt = ReactiveRuntime::init(Arc::new(|| {}));
let ui_thread = std::thread::current().id();
let observed_thread: Arc<Mutex<Option<std::thread::ThreadId>>> = Arc::new(Mutex::new(None));
let before_cross: Arc<Mutex<Option<ImageSource>>> = Arc::new(Mutex::new(None));
let bytes = two_by_two_png();
let owner = Owner::new();
let task = {
let observed_thread = observed_thread.clone();
let before_cross = before_cross.clone();
owner.with(|| {
use_task(move || {
let observed_thread = observed_thread.clone();
let before_cross = before_cross.clone();
let bytes = bytes.clone();
async move {
decode_with_probe(bytes, move |thread_id, result| {
*observed_thread.lock().expect("observed_thread poisoned") =
Some(thread_id);
if let Ok(source) = result {
*before_cross.lock().expect("before_cross poisoned") =
Some(source.clone());
}
})
.await
}
})
})
};
assert!(
pump_until(rt, Duration::from_secs(5), || task
.signal()
.get_untracked()
.is_ready()),
"decode_image_async must resolve to Ready via use_task"
);
let decoded_thread = observed_thread
.lock()
.expect("observed_thread poisoned")
.expect("probe must have recorded a thread id");
assert_ne!(
decoded_thread, ui_thread,
"the decode must run off the calling (UI) thread"
);
let after_cross = task
.signal()
.get_untracked()
.ready()
.cloned()
.expect("Ready must carry the decoded ImageSource");
let before_cross = before_cross
.lock()
.expect("before_cross poisoned")
.clone()
.expect("probe must have captured a pre-cross clone");
assert!(
before_cross.same(&after_cross),
"the ImageSource observed inside spawn_blocking (pre-cross) must share the \
same Arc allocation as the value observed on the UI thread (post-cross) — \
the result crosses as an Arc move, never a re-decode or pixel copy"
);
owner.cleanup();
}
#[test]
fn decode_image_async_composes_with_use_task() {
let rt = ReactiveRuntime::init(Arc::new(|| {}));
let bytes = two_by_two_png();
let owner = Owner::new();
let task = owner.with(|| {
use_task(move || {
let bytes = bytes.clone();
async move { decode_image_async(bytes).await }
})
});
assert!(
pump_until(rt, Duration::from_secs(5), || task
.signal()
.get_untracked()
.is_ready()),
"a valid PNG must decode to Ready"
);
let source = task.signal().get_untracked().ready().cloned().unwrap();
assert_eq!(source.natural_size(), kurbo::Size::new(2.0, 2.0));
owner.cleanup();
}
#[test]
fn decode_failure_surfaces_as_use_task_error() {
let rt = ReactiveRuntime::init(Arc::new(|| {}));
let owner = Owner::new();
let task = owner
.with(|| use_task(|| async { decode_image_async(b"not an image".to_vec()).await }));
assert!(
pump_until(rt, Duration::from_secs(5), || task
.signal()
.get_untracked()
.is_error()),
"invalid bytes must resolve to Error"
);
assert!(matches!(
task.signal().get_untracked().error().map(|e| e.to_string()),
Some(msg) if msg.contains("failed to decode image")
));
owner.cleanup();
}
}