#![warn(missing_docs)]
pub mod broadcaster;
pub mod ext;
pub mod operators;
pub mod sink;
pub mod source;
pub mod task;
pub use broadcaster::{BroadcastStream, Broadcaster, DEFAULT_BROADCAST_BUFFER};
pub use ext::RskitStreamExt;
pub use operators::combine::{concat, merge};
pub use sink::{collect, drain, for_each};
pub use source::{from_channel, from_fn, from_slice};
pub use task::{SpawnedTask, TaskGroup};
pub use tokio_util::sync::CancellationToken;
#[cfg(test)]
mod tests {
use parking_lot::Mutex;
use std::sync::Arc;
use std::time::Duration;
use futures::StreamExt as _;
use crate::{RskitStreamExt, from_fn, from_slice, merge};
#[tokio::test]
async fn test_from_slice_yields_all_in_order() {
let items = vec![1u32, 2, 3, 4, 5];
let stream = from_slice(items.clone());
let collected: Vec<u32> = stream.collect().await;
assert_eq!(collected, items);
}
#[tokio::test]
async fn test_from_slice_empty() {
let stream = from_slice::<u32>(vec![]);
let collected: Vec<u32> = stream.collect().await;
assert!(collected.is_empty());
}
#[tokio::test]
async fn test_from_fn_yields_until_none() {
let counter = Arc::new(Mutex::new(0u32));
let c = counter.clone();
let stream = from_fn(move || {
let c = c.clone();
async move {
let mut guard = c.lock();
let next = if *guard < 5 {
let val = *guard;
*guard += 1;
Some(val)
} else {
None
};
drop(guard);
next
}
});
let collected: Vec<u32> = stream.collect().await;
assert_eq!(collected, vec![0, 1, 2, 3, 4]);
}
#[tokio::test]
async fn test_from_fn_immediate_none() {
let stream = from_fn(|| async { None::<u32> });
let collected: Vec<u32> = stream.collect().await;
assert!(collected.is_empty());
}
#[tokio::test]
async fn test_merge_set_equality() {
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().await;
combined.sort_unstable();
assert_eq!(combined, vec![1, 2, 3, 4, 5, 6]);
}
#[tokio::test]
async fn test_merge_both_empty() {
let s1 = from_slice::<u32>(vec![]);
let s2 = from_slice::<u32>(vec![]);
let combined: Vec<u32> = merge(s1, s2).collect().await;
assert!(combined.is_empty());
}
#[tokio::test]
async fn test_rmap_transforms_items() {
let stream = from_slice(vec![1u32, 2, 3]);
let results: Vec<_> = stream
.rmap(|x| async move { Ok::<u32, rskit_errors::AppError>(x * 10) })
.collect()
.await;
let values: Vec<u32> = results.into_iter().map(|r| r.unwrap()).collect();
assert_eq!(values, vec![10, 20, 30]);
}
#[tokio::test]
async fn test_rmap_propagates_error() {
let stream = from_slice(vec![1u32, 2, 3]);
let results: Vec<_> = stream
.rmap(|x| async move {
if x == 2 {
Err(rskit_errors::AppError::new(
rskit_errors::ErrorCode::Internal,
"bad item",
))
} else {
Ok(x)
}
})
.collect()
.await;
assert!(results[0].is_ok());
assert!(results[1].is_err());
assert!(results[2].is_ok());
}
#[tokio::test]
async fn test_rfilter_keeps_matching_items() {
let stream = from_slice(vec![1u32, 2, 3, 4, 5, 6]);
let evens: Vec<u32> = stream.rfilter(|x| x % 2 == 0).collect().await;
assert_eq!(evens, vec![2, 4, 6]);
}
#[tokio::test]
async fn test_rfilter_no_match_yields_empty() {
let stream = from_slice(vec![1u32, 3, 5]);
let result: Vec<u32> = stream.rfilter(|x| x % 2 == 0).collect().await;
assert!(result.is_empty());
}
#[tokio::test]
async fn test_rtap_calls_side_effect_and_passes_through() {
let seen = Arc::new(Mutex::new(Vec::<u32>::new()));
let seen_clone = seen.clone();
let stream = from_slice(vec![10u32, 20, 30]);
let output: Vec<u32> = stream
.rtap(move |x| {
let seen = seen_clone.clone();
let val = *x;
async move {
seen.lock().push(val);
}
})
.collect()
.await;
assert_eq!(output, vec![10, 20, 30]);
assert_eq!(*seen.lock(), vec![10, 20, 30]);
}
#[tokio::test]
async fn test_rreduce_folds_to_single_value() {
let stream = from_slice(vec![1u32, 2, 3, 4, 5]);
let sum = stream.rreduce(0u32, |acc, x| acc + x).await;
assert_eq!(sum, 15);
}
#[tokio::test]
async fn test_rreduce_empty_stream_returns_init() {
let stream = from_slice::<u32>(vec![]);
let result = stream.rreduce(42u32, |acc, x| acc + x).await;
assert_eq!(result, 42);
}
#[tokio::test]
async fn test_rparallel_collects_all_results() {
let stream = from_slice(vec![1u32, 2, 3, 4, 5]);
let mut results: Vec<u32> = stream
.rparallel(
3,
|x| async move { Ok::<u32, rskit_errors::AppError>(x * 2) },
)
.collect::<Vec<_>>()
.await
.into_iter()
.map(|r| r.unwrap())
.collect();
results.sort_unstable();
assert_eq!(results, vec![2, 4, 6, 8, 10]);
}
#[tokio::test]
async fn test_rparallel_propagates_errors() {
let stream = from_slice(vec![1u32, 2, 3]);
let results: Vec<_> = stream
.rparallel(2, |x| async move {
if x == 2 {
Err(rskit_errors::AppError::new(
rskit_errors::ErrorCode::Internal,
"parallel error",
))
} else {
Ok(x)
}
})
.collect()
.await;
let error_count = results.iter().filter(|r| r.is_err()).count();
assert_eq!(error_count, 1);
}
#[tokio::test]
async fn test_rfan_out_applies_all_functions() {
let add_one = |x: u32| std::future::ready(Ok::<u32, rskit_errors::AppError>(x + 1));
let mul_two = |x: u32| std::future::ready(Ok::<u32, rskit_errors::AppError>(x * 2));
let stream_a = from_slice(vec![5u32, 10u32]);
let res_a: Vec<_> = stream_a.rfan_out(1, vec![add_one]).collect().await;
let res_a: Vec<Vec<_>> = res_a.into_iter().map(Result::unwrap).collect();
assert_eq!(res_a[0][0], 6u32);
assert_eq!(res_a[1][0], 11u32);
let stream_b = from_slice(vec![5u32, 10u32]);
let res_b: Vec<_> = stream_b.rfan_out(2, vec![add_one, mul_two]).collect().await;
let res_b: Vec<Vec<_>> = res_b.into_iter().map(Result::unwrap).collect();
assert_eq!(res_b[0][0], 6u32);
assert_eq!(res_b[0][1], 10u32);
assert_eq!(res_b[1][0], 11u32);
assert_eq!(res_b[1][1], 20u32);
}
#[tokio::test]
async fn test_rfan_out_single_function() {
let stream = from_slice(vec![3u32, 7u32]);
let f = |x: u32| std::future::ready(Ok::<u32, rskit_errors::AppError>(x + 100));
let results: Vec<_> = stream.rfan_out(1, vec![f]).collect().await;
let results: Vec<Vec<_>> = results.into_iter().map(Result::unwrap).collect();
assert_eq!(results.len(), 2);
assert_eq!(results[0][0], 103u32);
assert_eq!(results[1][0], 107u32);
}
#[tokio::test]
async fn test_rbatch_exact_size_batches() {
tokio::time::pause();
let stream = from_slice(vec![1u32, 2, 3, 4, 5, 6]);
let handle = tokio::spawn(async move {
stream
.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 test_rbatch_partial_flush_on_timeout() {
tokio::time::pause();
let (tx, rx) = tokio::sync::mpsc::channel::<u32>(16);
let stream = crate::source::from_channel(rx);
let handle = tokio::spawn(async move {
stream
.rbatch(10, Duration::from_millis(100))
.collect::<Vec<_>>()
.await
});
tx.send(1).await.unwrap();
tx.send(2).await.unwrap();
drop(tx);
tokio::time::advance(Duration::from_millis(200)).await;
let batches = handle.await.unwrap();
assert_eq!(batches.len(), 1);
assert_eq!(batches[0], vec![1, 2]);
}
#[tokio::test]
async fn test_rdebounce_emits_last_item() {
tokio::time::pause();
let (tx, rx) = tokio::sync::mpsc::channel::<u32>(16);
let stream = crate::source::from_channel(rx);
let handle = tokio::spawn(async move {
stream
.rdebounce(Duration::from_millis(100))
.collect::<Vec<_>>()
.await
});
tx.send(1).await.unwrap();
tx.send(2).await.unwrap();
tx.send(3).await.unwrap();
drop(tx);
tokio::time::advance(Duration::from_millis(200)).await;
let result = handle.await.unwrap();
assert!(!result.is_empty());
assert_eq!(*result.last().unwrap(), 3u32);
}
#[tokio::test]
async fn test_rthrottle_drops_fast_items() {
tokio::time::pause();
let stream = from_slice(vec![1u32, 2, 3, 4, 5]);
let handle = tokio::spawn(async move {
stream
.rthrottle(Duration::from_millis(100))
.collect::<Vec<_>>()
.await
});
tokio::time::advance(Duration::from_millis(600)).await;
let result = handle.await.unwrap();
assert!(!result.is_empty());
assert_eq!(result[0], 1u32);
assert!(result.len() < 5);
}
#[tokio::test]
async fn test_rtumbling_window_emits_on_timer() {
tokio::time::pause();
let (tx, rx) = tokio::sync::mpsc::channel::<u32>(16);
let stream = crate::source::from_channel(rx);
let handle = tokio::spawn(async move {
stream
.rtumbling_window(Duration::from_millis(100), 128)
.collect::<Vec<_>>()
.await
});
tx.send(10).await.unwrap();
tx.send(20).await.unwrap();
tx.send(30).await.unwrap();
drop(tx);
tokio::time::advance(Duration::from_millis(200)).await;
let windows = handle.await.unwrap();
assert!(!windows.is_empty());
let all_items: Vec<u32> = windows.into_iter().flatten().collect();
let mut sorted = all_items;
sorted.sort_unstable();
assert_eq!(sorted, vec![10, 20, 30]);
}
#[tokio::test]
async fn test_rtumbling_window_empty_input() {
tokio::time::pause();
let stream = from_slice::<u32>(vec![]);
let handle = tokio::spawn(async move {
stream
.rtumbling_window(Duration::from_millis(100), 128)
.collect::<Vec<_>>()
.await
});
tokio::time::advance(Duration::from_millis(200)).await;
let windows = handle.await.unwrap();
assert!(windows.is_empty());
}
#[tokio::test]
async fn test_rdistinct_filters_duplicates() {
let stream = from_slice(vec![1u32, 2, 2, 3, 1, 4]);
let values: Vec<u32> = stream.rdistinct().collect().await;
assert_eq!(values, vec![1, 2, 3, 4]);
}
#[tokio::test]
async fn test_rtake_and_rskip_compose() {
let stream = from_slice(vec![1u32, 2, 3, 4, 5]);
let values: Vec<u32> = stream.rskip(1).rtake(3).collect().await;
assert_eq!(values, vec![2, 3, 4]);
}
#[tokio::test]
async fn test_rpartition_splits_stream() {
let stream = from_slice(vec![1u32, 2, 3, 4, 5, 6]);
let (even_stream, odd_stream) = stream.rpartition(|value| value % 2 == 0);
let (evens, odds) = tokio::join!(
even_stream.collect::<Vec<_>>(),
odd_stream.collect::<Vec<_>>()
);
assert_eq!(evens, vec![2, 4, 6]);
assert_eq!(odds, vec![1, 3, 5]);
}
#[tokio::test]
async fn test_rsliding_window_emits_overlapping_windows() {
let stream = from_slice(vec![1u32, 2, 3, 4, 5]);
let windows: Vec<Vec<u32>> = stream.rsliding_window(3, 1).collect().await;
assert_eq!(windows, vec![vec![1, 2, 3], vec![2, 3, 4], vec![3, 4, 5]]);
}
}