use fluxion_stream::{MapOrderedExt, OrderedStreamExt, SkipItemsExt};
use fluxion_test_utils::helpers::unwrap_stream;
use fluxion_test_utils::test_data::{
person_alice, person_bob, person_charlie, person_dave, person_diane, TestData,
};
use fluxion_test_utils::{test_channel, unwrap_value, Sequenced};
#[tokio::test]
async fn test_map_ordered_skip_items() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream
.map_ordered(|item| Sequenced::new(item.into_inner())) .skip_items(2);
tx.unbounded_send(Sequenced::new(person_alice()))?; tx.unbounded_send(Sequenced::new(person_dave()))?; tx.unbounded_send(Sequenced::new(person_bob()))?; tx.unbounded_send(Sequenced::new(person_charlie()))?; tx.unbounded_send(Sequenced::new(person_diane()))?;
assert_eq!(
unwrap_stream(&mut result, 100).await.unwrap().into_inner(),
person_bob()
);
assert_eq!(
unwrap_stream(&mut result, 100).await.unwrap().into_inner(),
person_charlie()
);
assert_eq!(
unwrap_stream(&mut result, 100).await.unwrap().into_inner(),
person_diane()
);
Ok(())
}
#[tokio::test]
async fn test_ordered_merge_then_skip_items() -> anyhow::Result<()> {
let (s1_tx, s1_rx) = test_channel::<Sequenced<TestData>>();
let (s2_tx, s2_rx) = test_channel::<Sequenced<TestData>>();
let mut stream = s1_rx.ordered_merge(vec![s2_rx]).skip_items(3);
s1_tx.unbounded_send(Sequenced::new(person_alice()))?;
s2_tx.unbounded_send(Sequenced::new(person_bob()))?;
s1_tx.unbounded_send(Sequenced::new(person_charlie()))?;
s2_tx.unbounded_send(Sequenced::new(person_dave()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut stream, 500).await)).value,
person_dave()
);
s1_tx.unbounded_send(Sequenced::new(person_diane()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut stream, 500).await)).value,
person_diane()
);
Ok(())
}