use std::pin::pin;
use std::time::Duration;
use futures::StreamExt;
use ruststream::{
Broker, ConnectedBroker, Headers, IncomingMessage, OutgoingMessage, Publisher, Subscriber,
};
use ruststream_gcp_pubsub::{
ConnectedPubSubBroker, PARTITION_KEY_HEADER, PubSubBroker, PubSubSubscription,
};
const RECV_TIMEOUT: Duration = Duration::from_secs(15);
const TEST_PROJECT: &str = "ruststream-test";
fn test_host() -> Option<String> {
match std::env::var("PUBSUB_TEST_HOST") {
Ok(host) if !host.is_empty() => Some(host),
_ => {
eprintln!("PUBSUB_TEST_HOST is not set; skipping the emulator integration test");
None
}
}
}
async fn connect(host: &str) -> ConnectedPubSubBroker {
PubSubBroker::new(TEST_PROJECT)
.emulator(host)
.connect()
.await
.expect("broker connects")
}
fn unique(name: &str) -> String {
format!("it-{name}-{}", std::process::id())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn roundtrip_preserves_payload_attributes_and_partition_key() {
let Some(host) = test_host() else { return };
let connected = connect(&host).await;
let name = unique("roundtrip");
let mut subscriber = connected
.subscribe_descriptor(PubSubSubscription::new(&name).create_with_topic(&name))
.await
.expect("subscription opens");
let mut headers = Headers::new();
headers.insert("content-type", "application/json");
headers.insert("x-tenant", "acme");
headers.insert(PARTITION_KEY_HEADER, "user-42");
let publisher = connected.publisher();
publisher
.publish(OutgoingMessage::new(&name, b"{\"id\":1}".as_slice()).with_headers(headers))
.await
.expect("publish succeeds");
let mut stream = pin!(subscriber.stream());
let message = tokio::time::timeout(RECV_TIMEOUT, stream.next())
.await
.expect("delivery arrives")
.expect("stream is open")
.expect("delivery is ok");
assert_eq!(message.payload(), b"{\"id\":1}");
assert_eq!(
message.headers().get_str("content-type"),
Some("application/json")
);
assert_eq!(message.headers().get_str("x-tenant"), Some("acme"));
assert_eq!(message.partition_key(), Some(b"user-42".as_slice()));
message.ack().await.expect("ack succeeds");
connected.shutdown().await.expect("shutdown succeeds");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn nack_with_requeue_redelivers() {
let Some(host) = test_host() else { return };
let connected = connect(&host).await;
let name = unique("requeue");
let mut subscriber = connected
.subscribe_descriptor(PubSubSubscription::new(&name).create_with_topic(&name))
.await
.expect("subscription opens");
let publisher = connected.publisher();
publisher
.publish(OutgoingMessage::new(&name, b"again".as_slice()))
.await
.expect("publish succeeds");
let mut stream = pin!(subscriber.stream());
let first = tokio::time::timeout(RECV_TIMEOUT, stream.next())
.await
.expect("delivery arrives")
.expect("stream is open")
.expect("delivery is ok");
first.nack(true).await.expect("nack succeeds");
let second = tokio::time::timeout(RECV_TIMEOUT, stream.next())
.await
.expect("redelivery arrives")
.expect("stream is open")
.expect("redelivery is ok");
assert_eq!(second.payload(), b"again");
second.ack().await.expect("ack succeeds");
connected.shutdown().await.expect("shutdown succeeds");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn nack_without_requeue_does_not_redeliver() {
let Some(host) = test_host() else { return };
let connected = connect(&host).await;
let name = unique("drop");
let mut subscriber = connected
.subscribe_descriptor(PubSubSubscription::new(&name).create_with_topic(&name))
.await
.expect("subscription opens");
let publisher = connected.publisher();
publisher
.publish(OutgoingMessage::new(&name, b"poison".as_slice()))
.await
.expect("publish succeeds");
let mut stream = pin!(subscriber.stream());
let poison = tokio::time::timeout(RECV_TIMEOUT, stream.next())
.await
.expect("delivery arrives")
.expect("stream is open")
.expect("delivery is ok");
poison.nack(false).await.expect("drop succeeds");
publisher
.publish(OutgoingMessage::new(&name, b"next".as_slice()))
.await
.expect("publish succeeds");
let next = tokio::time::timeout(RECV_TIMEOUT, stream.next())
.await
.expect("delivery arrives")
.expect("stream is open")
.expect("delivery is ok");
assert_eq!(next.payload(), b"next");
next.ack().await.expect("ack succeeds");
connected.shutdown().await.expect("shutdown succeeds");
}