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::streaming::StreamFrame;
use velo::transports::tcp::TcpTransportBuilder;
use velo::{StreamConfig, TcpConfig, Velo};
async fn make_node() -> (Arc<Velo>, Registry) {
let registry = Registry::new();
let metrics = Arc::new(VeloMetrics::register(®istry).expect("metrics"));
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind");
let transport = Arc::new(
TcpTransportBuilder::new()
.from_listener(listener)
.expect("from_listener")
.build()
.expect("build tcp transport"),
);
let node = Velo::builder()
.add_transport(transport)
.stream_config(StreamConfig::Tcp(Some(TcpConfig {
bind_addr: std::net::Ipv4Addr::LOCALHOST.into(),
})))
.expect("stream_config")
.metrics(Arc::clone(&metrics))
.build()
.await
.expect("build velo");
(node, registry)
}
async fn make_pair() -> (Arc<Velo>, Arc<Velo>, Registry) {
let (consumer, _consumer_registry) = make_node().await;
let (producer, producer_registry) = make_node().await;
consumer
.register_peer(producer.peer_info())
.expect("register producer on consumer");
producer
.register_peer(consumer.peer_info())
.expect("register consumer on producer");
tokio::time::sleep(Duration::from_millis(200)).await;
(consumer, producer, producer_registry)
}
async fn collect_items(
anchor: &mut velo::streaming::StreamAnchor<u32>,
expected: usize,
) -> Vec<u32> {
let mut items = Vec::with_capacity(expected);
let drain = async {
while let Some(frame) = anchor.next().await {
match frame.expect("no stream error") {
StreamFrame::Item(v) => items.push(v),
StreamFrame::Finalized => break,
other => panic!("unexpected frame: {other:?}"),
}
}
items
};
tokio::time::timeout(Duration::from_secs(30), drain)
.await
.expect("timed out draining anchor")
}
#[tokio::test(flavor = "multi_thread")]
async fn burst_preserves_order_and_completeness() {
const N: u32 = 5_000;
let (consumer, producer, _reg) = make_pair().await;
let mut anchor = consumer.create_anchor::<u32>();
let handle = anchor.handle();
let sender = producer
.attach_anchor::<u32>(handle)
.await
.expect("remote attach over tcp-stream");
tokio::spawn(async move {
for i in 0..N {
sender.send(i).await.expect("send");
}
sender.finalize().expect("finalize");
});
let items = collect_items(&mut anchor, N as usize).await;
assert_eq!(items.len(), N as usize, "every frame must arrive");
assert!(
items.iter().copied().eq(0..N),
"coalescing must preserve send order"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn burst_coalesces_writes() {
const N: u32 = 20_000;
let (consumer, producer, producer_registry) = make_pair().await;
let mut anchor = consumer.create_anchor::<u32>();
let handle = anchor.handle();
let sender = producer
.attach_anchor::<u32>(handle)
.await
.expect("remote attach over tcp-stream");
tokio::spawn(async move {
for i in 0..N {
sender.send(i).await.expect("send");
}
sender.finalize().expect("finalize");
});
let items = collect_items(&mut anchor, N as usize).await;
assert_eq!(items.len(), N as usize);
let snap = MetricSnapshot::from_registry(&producer_registry);
let frames = snap.counter("velo_streaming_frames_written_total", &[]);
let flushes = snap.counter("velo_streaming_egress_flushes_total", &[]);
assert!(flushes > 0.0, "producer must have flushed batches");
assert!(
frames >= N as f64,
"expected at least {N} frames written, saw {frames}"
);
let ratio = frames / flushes;
eprintln!(
"coalescing ratio: {frames} frames / {flushes} flushes = {ratio:.2} frames per flush"
);
assert!(
ratio > 1.5,
"expected write coalescing under a back-to-back burst, but got \
{frames} frames in {flushes} flushes (ratio {ratio:.2}). A ratio near \
1.0 means every frame still went out on its own."
);
}
#[tokio::test(flavor = "multi_thread")]
async fn terminal_in_same_batch_flushes_preceding_frames() {
const N: u32 = 512;
let (consumer, producer, _reg) = make_pair().await;
let mut anchor = consumer.create_anchor::<u32>();
let handle = anchor.handle();
let sender = producer
.attach_anchor::<u32>(handle)
.await
.expect("remote attach over tcp-stream");
tokio::spawn(async move {
for i in 0..N {
sender.send(i).await.expect("send");
}
sender.finalize().expect("finalize");
});
let items = collect_items(&mut anchor, N as usize).await;
assert_eq!(
items.len(),
N as usize,
"frames staged ahead of the terminal must not be discarded"
);
assert!(items.iter().copied().eq(0..N), "order preserved");
}
#[tokio::test(flavor = "multi_thread")]
async fn forward_pass_shape_does_not_coalesce_per_stream() {
const STREAMS: usize = 32;
const PASSES: u32 = 100;
let (consumer, producer, producer_registry) = make_pair().await;
let mut anchors = Vec::with_capacity(STREAMS);
let mut senders = Vec::with_capacity(STREAMS);
for _ in 0..STREAMS {
let anchor = consumer.create_anchor::<u32>();
let handle = anchor.handle();
let sender = producer
.attach_anchor::<u32>(handle)
.await
.expect("remote attach");
anchors.push(anchor);
senders.push(sender);
}
let drive = tokio::spawn(async move {
for pass in 0..PASSES {
for sender in &senders {
sender.send(pass).await.expect("send");
}
tokio::time::sleep(Duration::from_millis(1)).await;
}
for sender in senders {
sender.finalize().expect("finalize");
}
});
for (idx, anchor) in anchors.iter_mut().enumerate() {
let items = collect_items(anchor, PASSES as usize).await;
assert_eq!(items.len(), PASSES as usize, "stream {idx} lost frames");
}
drive.await.expect("driver");
let snap = MetricSnapshot::from_registry(&producer_registry);
let frames = snap.counter("velo_streaming_frames_written_total", &[]);
let flushes = snap.counter("velo_streaming_egress_flushes_total", &[]);
let ratio = frames / flushes;
eprintln!(
"forward-pass shape: {frames} frames / {flushes} flushes = {ratio:.2} \
frames per flush across {STREAMS} streams"
);
assert!(
ratio < 1.5,
"expected per-stream coalescing to be ineffective for the \
one-frame-per-stream-per-pass shape, but got {ratio:.2}. If this is \
genuinely higher now, re-evaluate whether multiplexing is still needed."
);
}
#[tokio::test(flavor = "multi_thread")]
async fn concurrent_streams_stay_independent() {
const STREAMS: usize = 24;
const N: u32 = 200;
let (consumer, producer, _reg) = make_pair().await;
let mut anchors = Vec::with_capacity(STREAMS);
let mut senders = Vec::with_capacity(STREAMS);
for _ in 0..STREAMS {
let anchor = consumer.create_anchor::<u32>();
let handle = anchor.handle();
let sender = producer
.attach_anchor::<u32>(handle)
.await
.expect("remote attach");
anchors.push(anchor);
senders.push(sender);
}
for (idx, sender) in senders.into_iter().enumerate() {
let base = (idx as u32) * 1_000;
tokio::spawn(async move {
for i in 0..N {
sender.send(base + i).await.expect("send");
}
sender.finalize().expect("finalize");
});
}
for (idx, anchor) in anchors.iter_mut().enumerate() {
let base = (idx as u32) * 1_000;
let items = collect_items(anchor, N as usize).await;
assert_eq!(items.len(), N as usize, "stream {idx} lost frames");
assert!(
items.iter().copied().eq(base..base + N),
"stream {idx} received out-of-order or cross-delivered frames"
);
}
}