#![cfg(all(
feature = "memory",
feature = "macros",
feature = "json",
feature = "testing"
))]
mod common;
use std::time::Duration;
use ruststream::memory::ConnectedMemoryBroker;
use ruststream::memory::prelude::*;
use ruststream::testing::expect_published;
use common::{Event, Wire, connected};
#[subscriber("pairing.seeded")]
async fn consume(_event: &Event) -> HandlerOutcome {
HandlerOutcome::ack()
}
fn replaying() -> MemoryBroker<Retaining> {
MemoryBroker::retaining(Retention::Messages(nonzero!(8)))
}
async fn expect_payload(observer: &ConnectedMemoryBroker<Retaining>, name: &str, expected: &[u8]) {
let seen = expect_published(observer, name, 1, Duration::from_secs(2)).await;
assert_eq!(seen.len(), 1, "expected one publish on {name}");
assert_eq!(seen[0].payload(), expected);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_scope_hook_publishes_first_with_a_paired_publisher() {
let broker = replaying();
let observer = connected(&broker).await;
let app = RustStream::new(AppInfo::new("pairing", "0.1.0")).with_broker(broker, |b| {
b.include(consume);
b.after_startup(Publish, async move |publisher| {
publisher
.message(&Wire::of(b"first"))
.to("pairing.seeded")
.publish()
.await
});
});
let running = app.start().await.expect("startup failed");
expect_payload(&observer, "pairing.seeded", b"first").await;
running.shutdown().await.expect("graceful shutdown failed");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_running_handle_pairs_a_token_for_sibling_tasks() {
let broker = replaying().bindable();
let observer = connected(broker.broker()).await;
let egress = broker.bind(Publish);
let app = RustStream::new(AppInfo::new("pairing", "0.1.0")).with_broker(broker, |b| {
b.include(consume);
});
let running = app.start().await.expect("startup failed");
let publisher = running
.publisher(egress)
.await
.expect("pairing after start is infallible for memory");
publisher
.message(&Wire::of(b"late"))
.to("pairing.sibling")
.publish()
.await
.expect("publish");
expect_payload(&observer, "pairing.sibling", b"late").await;
running.shutdown().await.expect("graceful shutdown failed");
}
#[tokio::test]
async fn pairing_before_startup_reports_a_clear_error() {
let broker = MemoryBroker::new().bindable();
let token = broker.bind(Publish);
let _app = RustStream::new(AppInfo::new("pairing", "0.1.0")).with_broker(broker, |_b| {});
let err = token
.live()
.await
.expect_err("pairing before startup must fail");
assert!(err.to_string().contains("not connected"), "{err}");
}