tenshift-core 0.1.3

Thread-safe, backpressure-aware data loading pipeline for iterative processing
Documentation
use std::time::Instant;
use tenshift_core::sample::{Sample, Tensor};
use tenshift_core::sources::MemorySource;
use tenshift_core::Pipeline;

fn make_samples(n: usize) -> Vec<Sample> {
    (0..n)
        .map(|i| {
            // Simulate a real sample: 784 floats (28x28 image) + label
            Sample::new()
                .with("x", Tensor::f32(&vec![i as f32; 784], vec![28, 28]))
                .with("y", Tensor::i64(&[i as i64 % 10], vec![1]))
                .with_metadata(format!("sample_{i}"), i as u64)
        })
        .collect()
}

fn main() {
    println!("tenshift pipeline benchmark");
    println!("==============================\n");

    for dataset_size in [1_000, 10_000, 100_000] {
        let samples = make_samples(dataset_size);

        // Unbatched
        let start = Instant::now();
        let count: usize = Pipeline::from_source(MemorySource::new("bench", samples.clone()))
            .workers(1)
            .prefetch(4)
            .start()
            .unwrap()
            .map(|b| b.len())
            .sum();
        let single = start.elapsed();

        // Parallel workers
        let start = Instant::now();
        let count2: usize = Pipeline::from_source(MemorySource::new("bench", samples.clone()))
            .workers(4)
            .prefetch(8)
            .start()
            .unwrap()
            .map(|b| b.len())
            .sum();
        let parallel = start.elapsed();

        // With map transform (simulates CPU work)
        let start = Instant::now();
        let count3: usize = Pipeline::from_source(MemorySource::new("bench", samples.clone()))
            .workers(4)
            .map(|mut s| {
                // Simulate some CPU work: clone the tensor data
                let x = s.get("x").unwrap().try_as_f32().unwrap().to_vec();
                let doubled: Vec<f32> = x.iter().map(|v| v * 2.0).collect();
                s.insert("x", Tensor::f32(&doubled, vec![28, 28]));
                Ok(s)
            })
            .batch(32)
            .prefetch(8)
            .start()
            .unwrap()
            .map(|b| b.len())
            .sum();
        let with_work = start.elapsed();

        assert_eq!(count, dataset_size);
        assert_eq!(count2, dataset_size);
        assert_eq!(count3, dataset_size);

        let rate_s = dataset_size as f64 / single.as_secs_f64();
        let rate_p = dataset_size as f64 / parallel.as_secs_f64();
        let rate_w = dataset_size as f64 / with_work.as_secs_f64();

        println!("  {dataset_size} samples:");
        println!(
            "    1 worker passthrough: {:.1}ms ({:.0} samples/sec)",
            single.as_secs_f64() * 1000.0,
            rate_s,
        );
        println!(
            "    4 workers passthrough: {:.1}ms ({:.0} samples/sec)",
            parallel.as_secs_f64() * 1000.0,
            rate_p,
        );
        println!(
            "    4 workers + map + batch: {:.1}ms ({:.0} samples/sec)",
            with_work.as_secs_f64() * 1000.0,
            rate_w,
        );
        println!();
    }
}