use std::sync::Arc;
use std::time::Duration;
use futures::StreamExt;
use prometheus::Registry;
use velo::observability::VeloMetrics;
use velo::observability::test_helpers::MetricSnapshot;
use velo::transports::tcp::TcpTransportBuilder;
use velo::*;
fn new_messenger_transport() -> Arc<velo::transports::tcp::TcpTransport> {
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
Arc::new(
TcpTransportBuilder::new()
.from_listener(listener)
.unwrap()
.build()
.unwrap(),
)
}
async fn make_pair() -> (Arc<Velo>, Registry, Arc<Velo>, Registry) {
let producer_reg = Registry::new();
let producer_metrics = Arc::new(VeloMetrics::register(&producer_reg).unwrap());
let consumer_reg = Registry::new();
let consumer_metrics = Arc::new(VeloMetrics::register(&consumer_reg).unwrap());
let producer = Velo::builder()
.add_transport(new_messenger_transport())
.metrics(producer_metrics)
.build()
.await
.unwrap();
let consumer = Velo::builder()
.add_transport(new_messenger_transport())
.metrics(consumer_metrics)
.build()
.await
.unwrap();
producer.register_peer(consumer.peer_info()).unwrap();
consumer.register_peer(producer.peer_info()).unwrap();
tokio::time::sleep(Duration::from_millis(150)).await;
(producer, producer_reg, consumer, consumer_reg)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn reader_pump_backpressure_ticks_under_consumer_starvation() {
let (producer, _producer_reg, consumer, consumer_reg) = make_pair().await;
let mut anchor = consumer.create_anchor::<u64>();
let handle = anchor.handle();
let sender = producer
.attach_anchor::<u64>(handle)
.await
.expect("attach must succeed when both peers are registered");
let producer_task = tokio::spawn(async move {
for i in 0u64..8192 {
if sender.send(i).await.is_err() {
break;
}
}
tokio::time::sleep(Duration::from_millis(500)).await;
sender.finalize().ok();
});
let deadline = std::time::Instant::now() + Duration::from_secs(2);
let mut reader_pump_bp = 0.0;
let mut snap = MetricSnapshot::from_registry(&consumer_reg);
while std::time::Instant::now() < deadline {
snap = MetricSnapshot::from_registry(&consumer_reg);
reader_pump_bp = snap.counter("velo_streaming_reader_pump_backpressure_total", &[]);
if reader_pump_bp > 0.0 {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
reader_pump_bp > 0.0,
"reader_pump backpressure counter must tick when consumer can't keep up within 2s; got {reader_pump_bp}"
);
let watchdog = snap.counter("velo_streaming_heartbeat_watchdog_firings_total", &[]);
assert_eq!(
watchdog, 0.0,
"watchdog should not fire within 300ms (heartbeat default is 5s × 3 = 15s); got {watchdog}"
);
while let Some(frame) = anchor.next().await {
if matches!(frame, Ok(streaming::StreamFrame::Finalized) | Err(_)) {
break;
}
}
let _ = producer_task.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn no_backpressure_counters_in_steady_state() {
let (producer, producer_reg, consumer, consumer_reg) = make_pair().await;
let mut anchor = consumer.create_anchor::<u64>();
let handle = anchor.handle();
let sender = producer.attach_anchor::<u64>(handle).await.unwrap();
let consumer_task = tokio::spawn(async move {
let mut got = 0;
while let Some(frame) = anchor.next().await {
match frame.unwrap() {
streaming::StreamFrame::Item(_) => got += 1,
streaming::StreamFrame::Finalized => break,
_ => {}
}
}
got
});
for i in 0u64..16 {
sender.send(i).await.unwrap();
tokio::time::sleep(Duration::from_millis(2)).await;
}
sender.finalize().unwrap();
let got = tokio::time::timeout(Duration::from_secs(5), consumer_task)
.await
.unwrap()
.unwrap();
assert_eq!(got, 16);
for (label, reg) in [("producer", &producer_reg), ("consumer", &consumer_reg)] {
let snap = MetricSnapshot::from_registry(reg);
for metric in [
"velo_streaming_reader_pump_backpressure_total",
"velo_streaming_server_pump_backpressure_total",
"velo_streaming_producer_send_backpressure_total",
"velo_streaming_heartbeat_watchdog_firings_total",
] {
let v = snap.counter(metric, &[]);
assert_eq!(
v, 0.0,
"{label}: {metric} must stay 0 in steady state; got {v}"
);
}
}
}