Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
893bbab
Optimize TokenSpeed multimodal tensor transport
yechank-nvidia Jun 4, 2026
0dac842
Tighten multimodal transport validation
yechank-nvidia Jun 4, 2026
663b820
Rename TokenSpeed tensor storage wrapper
yechank-nvidia Jun 4, 2026
c3c8761
Clean up TokenSpeed SHM payloads on send failures
yechank-nvidia Jun 5, 2026
4876fa9
style(tokenspeed): format shm servicer path
yechank-nvidia Jun 9, 2026
0bdfdb5
fix(ci): satisfy clippy for multimodal transport
yechank-nvidia Jun 9, 2026
1d5e30f
perf(multimodal): optimize TokenSpeed SHM serialization
yechank-nvidia Jun 16, 2026
1959b60
feat(multimodal): auto TokenSpeed MM tensor transport + clearer env n…
yechank-nvidia Jun 17, 2026
58bfaa1
feat(multimodal): sweep orphaned SHM files + harden auto transport
yechank-nvidia Jun 17, 2026
2112066
feat(multimodal): Prometheus metrics for TokenSpeed MM tensor transport
yechank-nvidia Jun 17, 2026
9ad1be0
fix(multimodal): default Qwen VL image resize to bicubic
yechank-nvidia Jun 17, 2026
dc9dda1
fix(mm): place image/video placeholders before text in chat content
yechank-nvidia Jun 21, 2026
646e934
perf(mm): parallelize Pillow-exact BICUBIC resize (bit-identical)
yechank-nvidia Jun 21, 2026
1165b1d
feat(mm): libjpeg-turbo decode + Qwen3-VL preprocessing for PIL/HF pa…
yechank-nvidia Jun 21, 2026
45e949f
perf(mm): parallelize normalize + patchify (bit-identical)
yechank-nvidia Jun 21, 2026
a3981d1
fix(mm): address review — portable turbojpeg link, SAFETY docs, requi…
yechank-nvidia Jun 21, 2026
ac18f91
fix(grpc): address review — secure /dev/shm files + cleanup on assemb…
yechank-nvidia Jun 21, 2026
4c53004
perf(grpc): address review — chunk u16 SHM conversion on the contiguo…
yechank-nvidia Jun 21, 2026
020d7b7
fix(mm): load libturbojpeg at runtime (dlopen) instead of linking
yechank-nvidia Jun 21, 2026
e026d7a
refactor(grpc): address review — keep MM transport engine-neutral (sl…
yechank-nvidia Jun 22, 2026
5eecda6
fix(mm): satisfy clippy -D warnings (CI)
yechank-nvidia Jun 22, 2026
fdee8fc
style: cargo +nightly fmt
yechank-nvidia Jun 22, 2026
71194e0
fix(mm): address PR review — PIL-exact video resize, safer SHM locality
yechank-nvidia Jun 22, 2026
28fed55
fix(grpc): hoist audio_url placeholder in String content format
yechank-nvidia Jun 22, 2026
278cab7
refactor(grpc): drop unsafe as_slice_memory_order fallback; document …
yechank-nvidia Jun 22, 2026
6d7a601
refactor(grpc): drop self-introduced legacy MM transport env aliases
yechank-nvidia Jun 22, 2026
bc7da28
feat(grpc): verify shared /dev/shm via filesystem-identity handshake
yechank-nvidia Jun 22, 2026
1c766de
test(mm): guard PIL-exact video resize bit-identity
yechank-nvidia Jun 22, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions crates/grpc_client/proto/tokenspeed_scheduler.proto
Original file line number Diff line number Diff line change
Expand Up @@ -123,9 +123,11 @@ message TensorData {
oneof payload {
// Current path: raw little-endian bytes carried in the gRPC message.
bytes inline = 3;
// Same-host CPU shared memory path.
// Same-host CPU shared memory path. This is the preferred large-payload
// transport when SMG and TokenSpeed share /dev/shm.
ShmHandle shm = 4;
// Transport-specific remote descriptor (NIXL/RDMA/object-store/etc.).
// Cross-node or non-shared-memory transport descriptor. NIXL is the
// expected remote transport for distributed multimodal tensor payloads.
RemoteTensorHandle remote = 5;
}
}
Expand Down
1 change: 1 addition & 0 deletions crates/multimodal/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ hf-hub = { version = "0.5.0", default-features = false, features = ["tokio", "ru
bytes = { version = "1.12.0", features = ["serde"] }
fast_image_resize = { version = "6.0.0", features = ["image"] }
image = { version = "0.25.10", default-features = false, features = ["png", "jpeg", "gif", "bmp", "ico", "tiff", "webp"] }
libloading = "0.8"
ndarray = "0.17"
once_cell = "1.21.4"
opencv = { version = "0.98.2", default-features = false, features = ["clang-runtime", "imgproc", "videoio"], optional = true }
Expand Down
182 changes: 182 additions & 0 deletions crates/multimodal/src/jpeg_turbo.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,182 @@
//! Runtime (dlopen) binding to libjpeg-turbo's TurboJPEG API for JPEG decode.
//!
//! PIL/Pillow (and therefore vLLM) decode JPEGs with libjpeg-turbo using its
//! default options: accurate (islow) integer IDCT and "fancy" (bilinear) chroma
//! upsampling. The pure-Rust `image`/`zune-jpeg` decoder differs by a few levels
//! per pixel, which the vision encoder amplifies into a large embedding shift,
//! making TokenSpeed's multimodal accuracy diverge from vLLM. Decoding through
//! libjpeg-turbo with the same defaults makes SMG's pixel values match vLLM's.
//!
//! We load libturbojpeg at RUNTIME via `dlopen` rather than linking it, so the
//! crate (and every consumer — including the Go/Python bindings and CI builds
//! that don't ship libturbojpeg) compiles on any platform with no build script
//! and no link-time dependency. Where the shared library is present (the serving
//! image), decode goes through it for PIL parity; where it's absent,
//! `decode_jpeg_rgb` returns `None` and the caller falls back to the pure-Rust
//! decoder. Default flags (0) select accurate DCT + fancy upsampling, matching
//! Pillow.
//!
//! This module is the crate's only FFI surface, so it locally overrides the
//! workspace-wide `unsafe_code = "deny"` for the C bindings.
#![allow(unsafe_code)]

use std::{
os::raw::{c_int, c_uchar, c_ulong, c_void},
sync::OnceLock,
};

use image::{DynamicImage, RgbImage};
use libloading::{Library, Symbol};

type TjHandle = *mut c_void;
const TJPF_RGB: c_int = 0;

type TjInitDecompress = unsafe extern "C" fn() -> TjHandle;
type TjDecompressHeader3 = unsafe extern "C" fn(
TjHandle,
*const c_uchar,
c_ulong,
*mut c_int,
*mut c_int,
*mut c_int,
*mut c_int,
) -> c_int;
type TjDecompress2 = unsafe extern "C" fn(
TjHandle,
*const c_uchar,
c_ulong,
*mut c_uchar,
c_int,
c_int,
c_int,
c_int,
c_int,
) -> c_int;
type TjDestroy = unsafe extern "C" fn(TjHandle) -> c_int;

/// Resolved TurboJPEG entry points. Holds the loaded `Library` so the function
/// pointers stay valid for the process lifetime.
struct TurboJpeg {
_lib: Library,
init: TjInitDecompress,
header: TjDecompressHeader3,
decompress: TjDecompress2,
destroy: TjDestroy,
}

// The function pointers are plain C entry points with no shared mutable state;
// the library handle is kept alive for the process and never mutated.
unsafe impl Send for TurboJpeg {}
unsafe impl Sync for TurboJpeg {}

fn load_turbojpeg() -> Option<TurboJpeg> {
// Try the runtime soname first (shipped by the runtime package), then the
// dev symlink and common macOS names.
const CANDIDATES: &[&str] = &[
"libturbojpeg.so.0",
"libturbojpeg.so",
"libturbojpeg.0.dylib",
"libturbojpeg.dylib",
];
// SAFETY: loading a system shared library by name; we only resolve the four
// documented TurboJPEG symbols below and keep the handle for their lifetime.
let lib = CANDIDATES
.iter()
.find_map(|name| unsafe { Library::new(name) }.ok())?;
// SAFETY: each symbol is resolved against the just-loaded library with the
// signature documented by the TurboJPEG API. We copy the bare function
// pointers out (dropping the borrowing `Symbol`s) and keep `lib` alive in
// the returned struct, so the pointers remain valid.
let (init, header, decompress, destroy) = unsafe {
let init: Symbol<TjInitDecompress> = lib.get(b"tjInitDecompress\0").ok()?;
let header: Symbol<TjDecompressHeader3> = lib.get(b"tjDecompressHeader3\0").ok()?;
let decompress: Symbol<TjDecompress2> = lib.get(b"tjDecompress2\0").ok()?;
let destroy: Symbol<TjDestroy> = lib.get(b"tjDestroy\0").ok()?;
(*init, *header, *decompress, *destroy)
};
Some(TurboJpeg {
_lib: lib,
init,
header,
decompress,
destroy,
})
}

/// Process-wide cached TurboJPEG binding, or `None` if the library is absent.
fn turbojpeg() -> Option<&'static TurboJpeg> {
static TJ: OnceLock<Option<TurboJpeg>> = OnceLock::new();
TJ.get_or_init(load_turbojpeg).as_ref()
}

/// True if `bytes` start with the JPEG SOI marker.
pub fn is_jpeg(bytes: &[u8]) -> bool {
bytes.len() >= 3 && bytes[0] == 0xFF && bytes[1] == 0xD8 && bytes[2] == 0xFF
}

/// Decode a JPEG to an RGB8 `DynamicImage` via libjpeg-turbo (PIL-compatible
/// defaults). Returns `None` on any failure — including libturbojpeg not being
/// available at runtime — so the caller can fall back to the pure-Rust decoder.
pub fn decode_jpeg_rgb(bytes: &[u8]) -> Option<DynamicImage> {
if !is_jpeg(bytes) {
return None;
}
let tj = turbojpeg()?;
// SAFETY:
// - `tj.init` returns a handle that is null-checked before use; every early
// return below calls `tj.destroy(handle)` first, so the handle is freed
// exactly once and never used after destruction.
// - `bytes` is a live `&[u8]`; its ptr/len describe a valid immutable region
// for the duration of the calls (libjpeg-turbo only reads it).
// - `buf` is an owned `Vec<u8>` sized to exactly `w*h*3` (overflow-checked)
// from the header dimensions, and is the sole alias passed to the decoder;
// with `pitch=0` (=> `w*3`) and `TJPF_RGB` the decoder writes at most
// `w*h*3` bytes, so no out-of-bounds write occurs. The image is built only
// after `rc == 0` confirms a successful, complete write.
unsafe {
let handle = (tj.init)();
if handle.is_null() {
return None;
}
let (mut w, mut h, mut subsamp, mut colorspace) = (0_i32, 0_i32, 0_i32, 0_i32);
let hdr = (tj.header)(
handle,
bytes.as_ptr(),
bytes.len() as c_ulong,
&mut w,
&mut h,
&mut subsamp,
&mut colorspace,
);
if hdr != 0 || w <= 0 || h <= 0 {
(tj.destroy)(handle);
return None;
}
let (wu, hu) = (w as usize, h as usize);
// Guard against absurd dimensions before allocating.
let nbytes = match wu.checked_mul(hu).and_then(|p| p.checked_mul(3)) {
Some(n) => n,
None => {
(tj.destroy)(handle);
return None;
}
};
let mut buf = vec![0_u8; nbytes];
let rc = (tj.decompress)(
handle,
bytes.as_ptr(),
bytes.len() as c_ulong,
buf.as_mut_ptr(),
w,
0, // pitch = 0 -> width * pixelsize
h,
TJPF_RGB,
0, // default flags: accurate IDCT + fancy upsampling (matches Pillow)
);
(tj.destroy)(handle);
if rc != 0 {
return None;
}
RgbImage::from_raw(w as u32, h as u32, buf).map(DynamicImage::ImageRgb8)
}
}
1 change: 1 addition & 0 deletions crates/multimodal/src/lib.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
pub mod error;
pub mod hasher;
pub mod hub;
pub mod jpeg_turbo;
pub mod media;
pub mod registry;
pub mod tracker;
Expand Down
24 changes: 18 additions & 6 deletions crates/multimodal/src/media.rs
Original file line number Diff line number Diff line change
Expand Up @@ -329,12 +329,24 @@ impl MediaConnector {
) -> Result<Arc<ImageFrame>, MediaConnectorError> {
let hash = crate::hasher::hash_image(&bytes);

let cursor = std::io::Cursor::new(bytes.clone());
let reader = image::ImageReader::new(cursor).with_guessed_format()?;

let image = task::spawn_blocking(move || reader.decode())
.await
.map_err(MediaConnectorError::Blocking)??;
// Decode JPEGs through libjpeg-turbo (PIL-compatible defaults: accurate
// IDCT + fancy upsampling) so pixel values match vLLM bit-for-bit; the
// pure-Rust decoder diverges by a few levels, which the vision encoder
// amplifies into an embedding shift. Non-JPEG inputs and any turbojpeg
// failure fall back to the `image` crate.
let bytes_for_decode = bytes.clone();
let image = task::spawn_blocking(
move || -> Result<image::DynamicImage, MediaConnectorError> {
if let Some(img) = crate::jpeg_turbo::decode_jpeg_rgb(&bytes_for_decode) {
return Ok(img);
}
let cursor = std::io::Cursor::new(bytes_for_decode);
let reader = image::ImageReader::new(cursor).with_guessed_format()?;
Ok(reader.decode()?)
},
)
.await
.map_err(MediaConnectorError::Blocking)??;

Ok(Arc::new(ImageFrame::new(
image, bytes, detail, source, hash,
Expand Down
Loading
Loading