use super::*;
#[cfg(target_arch = "aarch64")]
pub(super) type FuseKernel = crate::resize_neon::Neon;
#[cfg(target_arch = "x86_64")]
pub(super) type FuseKernel = crate::resize_avx2::Avx2;
#[cfg(any(target_arch = "aarch64", target_arch = "x86_64"))]
pub(super) struct FuseChunks {
rx: std::sync::mpsc::Receiver<(Vec<u8>, usize)>,
recycle: std::sync::mpsc::Sender<Vec<u8>>,
row_bytes: usize,
}
#[cfg(any(target_arch = "aarch64", target_arch = "x86_64"))]
impl FuseChunks {
fn for_each_row(self, mut f: impl FnMut(&[u8]) -> Result<()>) -> Result<()> {
while let Ok((buf, rows)) = self.rx.recv() {
for r in 0..rows {
f(&buf[r * self.row_bytes..(r + 1) * self.row_bytes])?;
}
let _ = self.recycle.send(buf);
}
Ok(())
}
}
#[cfg(any(target_arch = "aarch64", target_arch = "x86_64"))]
fn fused_decode_loop<R: std::io::BufRead, T: Send>(
started: &mut mozjpeg::decompress::DecompressStarted<R>,
dec_w: usize,
dec_h: usize,
runway: usize,
worker: impl FnOnce(FuseChunks) -> Result<T> + Send,
) -> Result<Option<(f64, T)>> {
let row_bytes = dec_w * 3;
let chunk_rows = (64 * 1024 / row_bytes).clamp(1, dec_h);
let (chunk_tx, chunk_rx) = std::sync::mpsc::sync_channel::<(Vec<u8>, usize)>(runway);
let (recycle_tx, recycle_rx) = std::sync::mpsc::channel::<Vec<u8>>();
std::thread::scope(|sc| -> Result<Option<(f64, T)>> {
let chunks = FuseChunks {
rx: chunk_rx,
recycle: recycle_tx,
row_bytes,
};
let spawned = std::thread::Builder::new()
.name("oximg-fuse".into())
.spawn_scoped(sc, move || worker(chunks));
let Ok(worker) = spawned else {
return Ok(None);
};
let t_decode = std::time::Instant::now();
let mut worker_gone = false;
let decode_result = (|| -> Result<()> {
let mut remaining = dec_h;
while remaining > 0 {
let mut buf = recycle_rx.try_recv().unwrap_or_default();
let want = remaining.min(chunk_rows) * row_bytes;
if buf.len() < want {
buf.resize(want, 0);
}
let got = started
.read_scanlines_into(&mut buf[..want])
.context("decode failed")?
.len();
anyhow::ensure!(
got > 0 && got % row_bytes == 0,
"decoder returned a partial row"
);
let rows = got / row_bytes;
remaining -= rows;
if chunk_tx.send((buf, rows)).is_err() {
worker_gone = true;
return Err(anyhow::anyhow!("fuse worker exited early").context(ServerFault));
}
}
Ok(())
})();
let decode_ms = t_decode.elapsed().as_secs_f64() * 1e3;
drop(chunk_tx);
let worker_result = worker
.join()
.map_err(|_| anyhow::anyhow!("fuse worker panicked").context(ServerFault))?;
if worker_gone {
let value = worker_result?;
decode_result?; Ok(Some((decode_ms, value)))
} else {
decode_result?;
Ok(Some((decode_ms, worker_result?)))
}
})
}
#[cfg_attr(
not(any(target_arch = "aarch64", target_arch = "x86_64")),
allow(unused_variables)
)]
pub(super) fn fused_resize_encode<R: std::io::BufRead>(
started: &mut mozjpeg::decompress::DecompressStarted<R>,
dec_w: usize,
dec_h: usize,
dst_w: usize,
dst_h: usize,
quality: f32,
icc: Option<&[u8]>,
) -> Result<Option<(Vec<u8>, f64)>> {
#[cfg(not(any(target_arch = "aarch64", target_arch = "x86_64")))]
{
Ok(None)
}
#[cfg(any(target_arch = "aarch64", target_arch = "x86_64"))]
{
let Ok(mut resizer) =
crate::resize_kernel::StreamResize::<FuseKernel>::new(dec_w, dec_h, dst_w, dst_h, 3)
else {
return Ok(None);
};
let resizer = &mut resizer;
let out = fused_decode_loop(started, dec_w, dec_h, 2, move |chunks| {
let fwd = fwd_lut_f32();
let back = back_lut();
let mut row8 = vec![0u8; dst_w * 3];
let mut comp = jpegli::Compress::new(jpegli::ColorSpace::JCS_RGB);
comp.set_size(dst_w, dst_h);
comp.set_quality(quality);
if jpegli_progressive() {
comp.set_progressive_mode();
}
let mut enc = comp.start_compress(Vec::with_capacity(64 * 1024))?;
if let Some(icc) = icc {
for chunk in icc_app2_chunks(icc) {
enc.write_marker(jpegli::Marker::APP(2), &chunk);
}
}
chunks.for_each_row(|src| {
let mut enc_result = Ok(());
resizer.push_row_u8(src, fwd, |_, out| {
for (d, &v) in row8.iter_mut().zip(out) {
*d = back[v as usize];
}
if enc_result.is_ok() {
enc_result = enc.write_scanlines(&row8);
}
});
enc_result
.context("fused encode failed")
.context(ServerFault)
})?;
anyhow::ensure!(
resizer.rows_emitted() == dst_h,
"decode ended before the image was complete"
);
enc.finish()
.context("fused encode finish failed")
.context(ServerFault)
})?;
Ok(out.map(|(decode_ms, bytes)| (bytes, decode_ms)))
}
}
#[cfg_attr(
not(any(target_arch = "aarch64", target_arch = "x86_64")),
allow(unused_variables)
)]
#[allow(clippy::too_many_arguments)]
pub(super) fn fused_resize_pixels<R: std::io::BufRead, T: Send>(
started: &mut mozjpeg::decompress::DecompressStarted<R>,
dec_w: usize,
dec_h: usize,
dst_w: usize,
dst_h: usize,
out8: &mut [u8],
runway: usize,
side: impl FnOnce() -> Result<T> + Send,
) -> Result<Option<(f64, T)>> {
#[cfg(not(any(target_arch = "aarch64", target_arch = "x86_64")))]
{
let _ = side;
Ok(None)
}
#[cfg(any(target_arch = "aarch64", target_arch = "x86_64"))]
{
let Ok(mut resizer) =
crate::resize_kernel::StreamResize::<FuseKernel>::new(dec_w, dec_h, dst_w, dst_h, 3)
else {
return Ok(None);
};
let resizer = &mut resizer;
fused_decode_loop(started, dec_w, dec_h, runway, move |chunks| {
let side_value = side()?;
let fwd = fwd_lut_f32();
let back = back_lut();
chunks.for_each_row(|src| {
resizer.push_row_u8(src, fwd, |oy, out| {
for (d, &v) in out8[oy * dst_w * 3..(oy + 1) * dst_w * 3]
.iter_mut()
.zip(out)
{
*d = back[v as usize];
}
});
Ok(())
})?;
anyhow::ensure!(
resizer.rows_emitted() == dst_h,
"decode ended before the image was complete"
);
Ok(side_value)
})
}
}
#[cfg(feature = "avif")]
#[cfg_attr(
not(any(target_arch = "aarch64", target_arch = "x86_64")),
allow(unused_variables)
)]
#[allow(clippy::too_many_arguments)]
pub(super) fn fused_resize_yuv<R: std::io::BufRead>(
started: &mut mozjpeg::decompress::DecompressStarted<R>,
dec_w: usize,
dec_h: usize,
dst_w: usize,
dst_h: usize,
params: &crate::avif::AvifParams,
y_plane: &mut [u16],
cb_plane: &mut [u16],
cr_plane: &mut [u16],
) -> Result<Option<(f64, crate::avif::SvtSession)>> {
#[cfg(not(any(target_arch = "aarch64", target_arch = "x86_64")))]
{
Ok(None)
}
#[cfg(any(target_arch = "aarch64", target_arch = "x86_64"))]
{
let Ok(mut resizer) =
crate::resize_kernel::StreamResize::<FuseKernel>::new(dec_w, dec_h, dst_w, dst_h, 3)
else {
return Ok(None);
};
let resizer = &mut resizer;
let cw = dst_w.div_ceil(2);
fused_decode_loop(started, dec_w, dec_h, 4, move |chunks| {
let session = crate::avif::start_color_session(dst_w, dst_h, params)?;
let fwd = fwd_lut_f32();
let back = back_lut();
let mut row8 = vec![0u8; dst_w * 3];
let mut prev_row = vec![0u8; dst_w * 3];
chunks.for_each_row(|src| {
resizer.push_row_u8(src, fwd, |oy, out| {
for (d, &v) in row8.iter_mut().zip(out) {
*d = back[v as usize];
}
crate::avif::luma_rows(&row8, 3, &mut y_plane[oy * dst_w..][..dst_w]);
if oy % 2 == 1 {
let cy = oy / 2;
crate::avif::chroma_row_pair(
&prev_row,
Some(&row8),
dst_w,
3,
&mut cb_plane[cy * cw..][..cw],
&mut cr_plane[cy * cw..][..cw],
);
} else {
prev_row.copy_from_slice(&row8);
}
});
Ok(())
})?;
anyhow::ensure!(
resizer.rows_emitted() == dst_h,
"decode ended before the image was complete"
);
if dst_h % 2 == 1 {
let cy = dst_h / 2;
crate::avif::chroma_row_pair(
&prev_row,
None,
dst_w,
3,
&mut cb_plane[cy * cw..][..cw],
&mut cr_plane[cy * cw..][..cw],
);
}
Ok(session)
})
}
}