use fluxion_core::CancellationToken;
use fluxion_exec::subscribe::SubscribeExt;
use futures::channel::mpsc::unbounded;
use futures::lock::Mutex as FutureMutex;
use std::sync::Arc;
use tokio::spawn;
use tokio_stream::StreamExt as _;
#[tokio::test]
async fn test_subscribe_example() -> anyhow::Result<()> {
#[derive(Debug, Clone, PartialEq, Eq)]
struct Item {
id: u32,
value: String,
}
#[derive(Debug, thiserror::Error)]
#[error("Test error: {0}")]
struct TestError(String);
let (tx, rx) = unbounded::<Item>();
let stream = rx;
let results = Arc::new(FutureMutex::new(Vec::new()));
let (notify_tx, mut notify_rx) = unbounded();
let process_func = {
let results = results.clone();
let notify_tx = notify_tx.clone();
move |item: Item, _ctx: CancellationToken| {
let results = results.clone();
let notify_tx = notify_tx.clone();
async move {
results.lock().await.push(item);
let _ = notify_tx.unbounded_send(()); Ok::<(), TestError>(())
}
}
};
let task = spawn(async move {
stream
.subscribe(process_func, |_| {}, None)
.await
.expect("subscribe should succeed");
});
let item1 = Item {
id: 1,
value: "Alice".to_string(),
};
let item2 = Item {
id: 2,
value: "Bob".to_string(),
};
let item3 = Item {
id: 3,
value: "Charlie".to_string(),
};
tx.unbounded_send(item1.clone())?;
notify_rx.next().await.unwrap();
assert_eq!(*results.lock().await, vec![item1.clone()]);
tx.unbounded_send(item2.clone())?;
notify_rx.next().await.unwrap();
assert_eq!(*results.lock().await, vec![item1.clone(), item2.clone()]);
tx.unbounded_send(item3.clone())?;
notify_rx.next().await.unwrap();
assert_eq!(*results.lock().await, vec![item1, item2, item3]);
drop(tx);
task.await.unwrap();
Ok(())
}