use fluxion_stream::combine_with_previous::CombineWithPreviousExt;
use fluxion_stream::prelude::*;
use fluxion_test_utils::test_data::{
person_alice, person_bob, person_charlie, person_dave, TestData,
};
use fluxion_test_utils::Sequenced;
use fluxion_test_utils::{assert_stream_ended, test_channel};
use fluxion_test_utils::{helpers::unwrap_stream, unwrap_value};
#[tokio::test]
async fn test_map_ordered_basic_transformation() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream.combine_with_previous().map_ordered(|stream_item| {
Sequenced::new(format!(
"Previous: {:?}, Current: {}",
stream_item.previous.map(|p| p.value.to_string()),
&stream_item.current.value
))
});
tx.unbounded_send(Sequenced::new(person_alice()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Previous: None, Current: Person[name=Alice, age=25]"
);
tx.unbounded_send(Sequenced::new(person_bob()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Previous: Some(\"Person[name=Alice, age=25]\"), Current: Person[name=Bob, age=30]"
);
tx.unbounded_send(Sequenced::new(person_charlie()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Previous: Some(\"Person[name=Bob, age=30]\"), Current: Person[name=Charlie, age=35]"
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_to_struct() -> anyhow::Result<()> {
#[derive(Debug, PartialEq, Clone, Ord, PartialOrd, Eq)]
struct AgeComparison {
previous_age: Option<u32>,
current_age: u32,
age_increased: bool,
}
impl AgeComparison {
fn new(previous_age: Option<u32>, current_age: u32) -> Self {
let age_increased = previous_age.is_some_and(|prev| current_age > prev);
AgeComparison {
previous_age,
current_age,
age_increased,
}
}
}
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream.combine_with_previous().map_ordered(|stream_item| {
let current_age = match &stream_item.current.value {
TestData::Person(p) => p.age,
_ => 0,
};
let previous_age = stream_item
.previous
.as_ref()
.and_then(|prev| match &prev.value {
TestData::Person(p) => Some(p.age),
_ => None,
});
Sequenced::new(AgeComparison::new(previous_age, current_age))
});
tx.unbounded_send(Sequenced::new(person_alice()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
AgeComparison::new(None, 25)
);
tx.unbounded_send(Sequenced::new(person_bob()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
AgeComparison::new(Some(25), 30)
);
tx.unbounded_send(Sequenced::new(person_alice()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
AgeComparison::new(Some(30), 25)
);
tx.unbounded_send(Sequenced::new(person_charlie()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
AgeComparison::new(Some(25), 35)
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_extract_age_difference() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream
.combine_with_previous()
.map_ordered(|stream_item| -> Sequenced<i32> {
let current_age = match &stream_item.current.value {
TestData::Person(p) => p.age as i32,
_ => 0,
};
let previous_age = stream_item
.previous
.as_ref()
.and_then(|prev| match &prev.value {
TestData::Person(p) => Some(p.age as i32),
_ => None,
});
Sequenced::new(current_age - previous_age.unwrap_or(current_age))
});
tx.unbounded_send(Sequenced::new(person_alice()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
0
);
tx.unbounded_send(Sequenced::new(person_bob()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
5
);
tx.unbounded_send(Sequenced::new(person_dave()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
-2
);
tx.unbounded_send(Sequenced::new(person_charlie()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
7
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_single_value() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream
.combine_with_previous()
.map_ordered(|stream_item| Sequenced::new(stream_item.current.value.to_string()));
tx.unbounded_send(Sequenced::new(person_alice()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Person[name=Alice, age=25]"
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_empty_stream() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream
.combine_with_previous()
.map_ordered(|stream_item| Sequenced::new(stream_item.current.value.to_string()));
drop(tx);
assert_stream_ended(&mut result, 500).await;
Ok(())
}
#[tokio::test]
async fn test_map_ordered_preserves_ordering() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream.combine_with_previous().map_ordered(|stream_item| {
Sequenced::new(match &stream_item.current.value {
TestData::Person(p) => p.name.clone(),
_ => String::from("Unknown"),
})
});
tx.unbounded_send(Sequenced::new(person_alice()))?;
tx.unbounded_send(Sequenced::new(person_bob()))?;
tx.unbounded_send(Sequenced::new(person_charlie()))?;
tx.unbounded_send(Sequenced::new(person_dave()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Alice"
);
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Bob"
);
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Charlie"
);
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
"Dave"
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_multiple_transformations() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream.combine_with_previous().map_ordered(|stream_item| {
Sequenced::new(match &stream_item.current.value {
TestData::Person(p) => p.age,
_ => 0,
})
});
tx.unbounded_send(Sequenced::new(person_alice()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
25
);
tx.unbounded_send(Sequenced::new(person_bob()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
30
);
tx.unbounded_send(Sequenced::new(person_charlie()))?;
assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
35
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_with_complex_closure() -> anyhow::Result<()> {
#[derive(Debug, PartialEq, Clone, Ord, PartialOrd, Eq)]
struct PersonSummary {
name: String,
age_category: &'static str,
changed_from_previous: bool,
}
impl PersonSummary {
fn new(name: String, age_category: &'static str, changed_from_previous: bool) -> Self {
PersonSummary {
name,
age_category,
changed_from_previous,
}
}
}
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream.combine_with_previous().map_ordered(|stream_item| {
let current = match &stream_item.current.value {
TestData::Person(p) => p,
_ => panic!("Expected person"),
};
let age_category = match current.age {
0..=17 => "child",
18..=29 => "young adult",
30..=59 => "adult",
_ => "senior",
};
let changed_from_previous = !stream_item.previous.as_ref().is_some_and(|prev| {
if let TestData::Person(prev_person) = &prev.value {
prev_person.name == current.name
} else {
false
}
});
Sequenced::new(PersonSummary::new(
current.name.clone(),
age_category,
changed_from_previous,
))
});
tx.unbounded_send(Sequenced::new(person_alice()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
PersonSummary::new(String::from("Alice"), "young adult", true)
);
tx.unbounded_send(Sequenced::new(person_bob()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
PersonSummary::new(String::from("Bob"), "adult", true)
);
tx.unbounded_send(Sequenced::new(person_bob()))?; assert_eq!(
unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value,
PersonSummary::new(String::from("Bob"), "adult", false)
);
Ok(())
}
#[tokio::test]
async fn test_map_ordered_boolean_logic() -> anyhow::Result<()> {
let (tx, stream) = test_channel::<Sequenced<TestData>>();
let mut result = stream.combine_with_previous().map_ordered(|stream_item| {
let current_age = match &stream_item.current.value {
TestData::Person(p) => p.age,
_ => 0,
};
Sequenced::new(stream_item.previous.as_ref().is_some_and(|prev| {
if let TestData::Person(p) = &prev.value {
current_age > p.age
} else {
false
}
}))
});
tx.unbounded_send(Sequenced::new(person_alice()))?; assert!(!unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value);
tx.unbounded_send(Sequenced::new(person_bob()))?; assert!(unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value);
tx.unbounded_send(Sequenced::new(person_charlie()))?; assert!(unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value);
tx.unbounded_send(Sequenced::new(person_dave()))?; assert!(!unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value);
tx.unbounded_send(Sequenced::new(person_alice()))?; assert!(!unwrap_value(Some(unwrap_stream(&mut result, 500).await)).value);
Ok(())
}