use fluxion_core::{FluxionError, StreamItem};
use fluxion_stream::EmitWhenExt;
use fluxion_test_utils::{test_channel_with_errors, unwrap_stream, Sequenced};
use futures::StreamExt;
#[tokio::test]
async fn test_emit_when_propagates_source_error() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(15, 1)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
source_tx.unbounded_send(StreamItem::Error(FluxionError::stream_error(
"Source error",
)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Error(_)
));
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(30, 3)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_propagates_filter_error() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 1)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
filter_tx.unbounded_send(StreamItem::Error(FluxionError::stream_error(
"Filter error",
)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Error(_)
));
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(15, 4)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(20, 3)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_predicate_continues_after_error() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(15, 1)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
source_tx.unbounded_send(StreamItem::Error(FluxionError::stream_error("Error")))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Error(_)
));
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(25, 4)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(30, 5)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_both_streams_have_errors() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 1)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
source_tx.unbounded_send(StreamItem::Error(FluxionError::stream_error(
"Source error",
)))?;
let result2 = result.next().await.unwrap();
assert!(matches!(result2, StreamItem::Error(_)));
filter_tx.unbounded_send(StreamItem::Error(FluxionError::stream_error(
"Filter error",
)))?;
let result3 = result.next().await.unwrap();
assert!(matches!(result3, StreamItem::Error(_)));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_error_before_filter_ready() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
source_tx.unbounded_send(StreamItem::Error(FluxionError::stream_error("Early error")))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Error(_)
));
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 2)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 1)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_source_none_on_filter_update() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 1)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_filter_returns_false_on_source_update() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(50, 1)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(100, 3)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_filter_returns_false_on_filter_update() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 1)))?;
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 2)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(20, 3)))?;
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(8, 4)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_source_none_on_source_update() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 1)))?;
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 2)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_multiple_filter_updates_no_source() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 1)))?;
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 2)))?;
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(3, 3)))?;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(15, 4)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}
#[tokio::test]
async fn test_emit_when_alternating_false_conditions() -> anyhow::Result<()> {
let (source_tx, source_stream) = test_channel_with_errors::<Sequenced<i32>>();
let (filter_tx, filter_stream) = test_channel_with_errors::<Sequenced<i32>>();
let mut result =
source_stream.emit_when(filter_stream, |state| state.values()[0] > state.values()[1]);
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(10, 1)))?;
filter_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(5, 2)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(3, 3)))?;
tokio::time::sleep(tokio::time::Duration::from_millis(50)).await;
source_tx.unbounded_send(StreamItem::Value(Sequenced::with_timestamp(20, 4)))?;
assert!(matches!(
unwrap_stream(&mut result, 100).await,
StreamItem::Value(_)
));
drop(source_tx);
drop(filter_tx);
Ok(())
}