use parking_lot::Mutex;
use std::sync::Arc;
use std::time::Duration;
use futures_util::StreamExt;
use rskit_errors::{AppError, AppResult, ErrorCode};
use rskit_stream::{RskitStreamExt, concat, from_slice, merge};
#[tokio::test]
async fn empty_stream_rmap_emits_no_items() {
let results: Vec<AppResult<u32>> = from_slice::<u32>(vec![])
.rmap(|x| async move { Ok(x * 10) })
.collect::<Vec<_>>()
.await;
assert!(results.is_empty());
}
#[tokio::test]
async fn empty_stream_rbatch_emits_no_batches() {
tokio::time::pause();
let handle = tokio::spawn(async {
from_slice::<u32>(vec![])
.rbatch(5, Duration::from_millis(100))
.collect::<Vec<_>>()
.await
});
tokio::time::advance(Duration::from_millis(200)).await;
let batches = handle.await.unwrap();
assert!(batches.is_empty());
}
#[tokio::test]
async fn single_item_pipeline_processes_once() {
let results: Vec<AppResult<u32>> = from_slice(vec![42u32])
.rmap(|x| async move { Ok(x + 1) })
.rfilter(|r| r.as_ref().is_ok_and(|v| *v > 0))
.collect::<Vec<_>>()
.await;
assert_eq!(results.len(), 1);
assert_eq!(results[0].as_ref().unwrap(), &43u32);
}
#[tokio::test]
async fn large_stream_rmap_processes_all_items() {
let items: Vec<u32> = (0..10_000).collect();
let results: Vec<AppResult<u32>> = from_slice(items)
.rmap(|x| async move { Ok(x * 2) })
.collect::<Vec<_>>()
.await;
assert_eq!(results.len(), 10_000);
for (i, r) in results.iter().enumerate() {
assert_eq!(*r.as_ref().unwrap(), u32::try_from(i).unwrap() * 2);
}
}
#[tokio::test]
async fn large_stream_rparallel_processes_all_items() {
let items: Vec<u32> = (0..10_000).collect();
let mut results: Vec<u32> = from_slice(items)
.rparallel(8, |x| async move { Ok::<u32, AppError>(x * 3) })
.collect::<Vec<_>>()
.await
.into_iter()
.map(|r| r.unwrap())
.collect();
results.sort_unstable();
let expected: Vec<u32> = (0..10_000).map(|x| x * 3).collect();
assert_eq!(results, expected);
}
#[tokio::test]
async fn rmap_error_recovery_continues_stream() {
let results: Vec<AppResult<u32>> = from_slice(vec![1u32, 2, 3, 4, 5])
.rmap(|x| async move {
if x % 2 == 0 {
Err(AppError::new(ErrorCode::Internal, "even number"))
} else {
Ok(x * 10)
}
})
.collect::<Vec<_>>()
.await;
assert_eq!(results.len(), 5);
assert_eq!(*results[0].as_ref().unwrap(), 10);
assert!(results[1].is_err());
assert_eq!(*results[2].as_ref().unwrap(), 30);
assert!(results[3].is_err());
assert_eq!(*results[4].as_ref().unwrap(), 50);
}
#[tokio::test]
async fn rparallel_surfaces_worker_errors() {
let results: Vec<AppResult<u32>> = from_slice(vec![1u32, 2, 3, 4, 5])
.rparallel(4, |x| async move {
if x == 3 || x == 5 {
Err(AppError::new(ErrorCode::Internal, "bad value"))
} else {
Ok(x)
}
})
.collect::<Vec<_>>()
.await;
assert_eq!(results.len(), 5);
let ok_count = results.iter().filter(|r| r.is_ok()).count();
let err_count = results.iter().filter(|r| r.is_err()).count();
assert_eq!(ok_count, 3);
assert_eq!(err_count, 2);
}
#[tokio::test]
async fn rfan_out_propagates_branch_errors() {
let ok_fn = |x: u32| std::future::ready(Ok::<u32, AppError>(x + 1));
let err_fn = |_x: u32| {
std::future::ready(Err::<u32, AppError>(AppError::new(
ErrorCode::Internal,
"fan_out error",
)))
};
let results: Vec<AppResult<Vec<u32>>> = from_slice(vec![10u32, 20])
.rfan_out(2, vec![ok_fn, err_fn])
.collect::<Vec<_>>()
.await;
assert_eq!(results.len(), 2);
assert!(results[0].is_err());
assert!(results[1].is_err());
}
#[tokio::test]
async fn complex_chain_of_five_operators_preserves_expected_items() {
let seen = Arc::new(Mutex::new(Vec::<u32>::new()));
let seen_clone = seen.clone();
let sum = from_slice(vec![1u32, 2, 3, 4, 5, 6, 7, 8, 9, 10])
.rmap(|x| async move { Ok::<u32, AppError>(x * 2) })
.rfilter(|r| r.as_ref().is_ok_and(|v| *v > 6))
.rtap(move |r| {
let seen = seen_clone.clone();
let val = r.as_ref().ok().copied();
async move {
if let Some(v) = val {
seen.lock().push(v);
}
}
})
.rreduce(0u32, |acc, r| acc + r.unwrap_or(0))
.await;
assert_eq!(sum, 98);
let tapped = seen.lock().clone();
assert_eq!(tapped, vec![8, 10, 12, 14, 16, 18, 20]);
}
#[tokio::test]
async fn rmap_then_rbatch_batches_mapped_values() {
tokio::time::pause();
let handle = tokio::spawn(async {
from_slice(vec![1u32, 2, 3, 4, 5])
.rmap(|x| async move { Ok::<u32, AppError>(x + 100) })
.rbatch(2, Duration::from_secs(10))
.collect::<Vec<_>>()
.await
});
tokio::time::advance(Duration::from_secs(11)).await;
let batches = handle.await.unwrap();
assert_eq!(batches.len(), 3);
assert_eq!(batches[0].len(), 2);
assert_eq!(batches[1].len(), 2);
assert_eq!(batches[2].len(), 1);
let all: Vec<u32> = batches
.into_iter()
.flatten()
.map(|r: AppResult<u32>| r.unwrap())
.collect();
assert_eq!(all, vec![101, 102, 103, 104, 105]);
}
#[tokio::test]
async fn merge_combines_multiple_streams() {
let s1 = from_slice(vec![1u32, 3, 5]);
let s2 = from_slice(vec![2u32, 4, 6]);
let mut combined: Vec<u32> = merge(s1, s2).collect::<Vec<_>>().await;
combined.sort_unstable();
assert_eq!(combined, vec![1, 2, 3, 4, 5, 6]);
}
#[tokio::test]
async fn concat_preserves_stream_order() {
let s1 = from_slice(vec![1u32, 2, 3]);
let s2 = from_slice(vec![4u32, 5, 6]);
let s3 = from_slice(vec![7u32, 8, 9]);
let result: Vec<u32> = concat(vec![s1, s2, s3]).collect::<Vec<_>>().await;
assert_eq!(result, vec![1, 2, 3, 4, 5, 6, 7, 8, 9]);
}
#[tokio::test]
async fn concat_skips_empty_streams() {
let s1 = from_slice::<u32>(vec![]);
let s2 = from_slice(vec![10u32, 20]);
let s3 = from_slice::<u32>(vec![]);
let s4 = from_slice(vec![30u32]);
let result: Vec<u32> = concat(vec![s1, s2, s3, s4]).collect::<Vec<_>>().await;
assert_eq!(result, vec![10, 20, 30]);
}
#[tokio::test]
async fn rbatch_exact_multiple_emits_full_batches() {
tokio::time::pause();
let handle = tokio::spawn(async {
from_slice(vec![1u32, 2, 3, 4, 5, 6])
.rbatch(3, Duration::from_millis(500))
.collect::<Vec<_>>()
.await
});
tokio::time::advance(Duration::from_millis(600)).await;
let batches = handle.await.unwrap();
assert_eq!(batches.len(), 2);
assert_eq!(batches[0], vec![1, 2, 3]);
assert_eq!(batches[1], vec![4, 5, 6]);
}
#[tokio::test]
async fn rbatch_single_item_batches_emit_each_item() {
tokio::time::pause();
let handle = tokio::spawn(async {
from_slice(vec![10u32, 20, 30])
.rbatch(1, Duration::from_millis(500))
.collect::<Vec<_>>()
.await
});
tokio::time::advance(Duration::from_millis(600)).await;
let batches = handle.await.unwrap();
assert_eq!(batches.len(), 3);
assert_eq!(batches[0], vec![10]);
assert_eq!(batches[1], vec![20]);
assert_eq!(batches[2], vec![30]);
}
#[tokio::test]
async fn rfilter_then_rparallel_processes_filtered_values() {
let mut results: Vec<u32> = from_slice(vec![1u32, 2, 3, 4, 5, 6, 7, 8, 9, 10])
.rfilter(|&x| x % 3 == 0)
.rparallel(4, |x| async move { Ok::<u32, AppError>(x * 100) })
.collect::<Vec<_>>()
.await
.into_iter()
.map(|r| r.unwrap())
.collect();
results.sort_unstable();
assert_eq!(results, vec![300, 600, 900]);
}
#[tokio::test]
async fn pipeline_composition_can_be_reused() {
let data = vec![1u32, 2, 3, 4, 5];
let chain_a: Vec<u32> = from_slice(data.clone())
.rmap(|x| async move { Ok::<u32, AppError>(x * 2) })
.collect::<Vec<_>>()
.await
.into_iter()
.map(|r| r.unwrap())
.collect();
assert_eq!(chain_a, vec![2, 4, 6, 8, 10]);
let chain_b = from_slice(data.clone())
.rfilter(|&x| x > 2)
.rreduce(0u32, |acc, x| acc + x)
.await;
assert_eq!(chain_b, 3 + 4 + 5);
let mut chain_c: Vec<u32> = from_slice(data)
.rmap(|x| async move { Ok::<u32, AppError>(x + 10) })
.rparallel(2, |r| async move {
let val = r?;
Ok(val * 2)
})
.collect::<Vec<_>>()
.await
.into_iter()
.map(|r| r.unwrap())
.collect();
chain_c.sort_unstable();
assert_eq!(chain_c, vec![22, 24, 26, 28, 30]);
}