#include <xpp/sync/mpsc.h>
xpp::EventLoop loop;
xpp::WaitScope scope(loop);
auto [tx, rx] = xpp::sync::mpsc::channel<int>(16);
// Async: .await() suspends (fiber) or blocks + drives loop
tx.send(42).await();
// Sync: returns immediately, fails when full
auto r = tx.try_send(99);
if (r.is_err()) { /* Full or Closed */ }
// Clone the sender for multiple producers
auto tx2 = tx;
tx2.send(10).await();
// Receive
auto v = rx.recv().await(); // Option<T> — none if closed
auto v = rx.try_recv(); // Result<T, TryRecvError>
With fiber — non-blocking:
xpp::fiber([]() {
auto [tx, rx] = xpp::sync::mpsc::channel<int>(4);
std::thread producer([tx = std::move(tx)]() mutable {
for (int i = 0; i < 10; i++) tx.send(i).await(); // blocks when full
}).detach();
for (int i = 0; i < 10; i++) {
auto v = rx.recv().await(); // fiber suspends when empty
printf("got %d\n", v.unwrap());
}
}).await();