mod common;
use common::MockFrameTransport;
use futures::StreamExt;
use std::sync::Arc;
use std::time::Duration;
use velo::streaming::{AnchorManager, AttachError, StreamError, StreamFrame};
use velo_ext::WorkerId;
fn make_manager() -> Arc<AnchorManager> {
Arc::new(AnchorManager::new(
WorkerId::from_u64(1),
Arc::new(MockFrameTransport::new()),
))
}
#[tokio::test(flavor = "multi_thread")]
async fn test_03_exclusive_reject() {
let mgr = make_manager();
let anchor = mgr.create_anchor::<u32>();
let handle = anchor.handle();
let _sender = mgr
.attach_stream_anchor::<u32>(handle)
.await
.expect("first attach must succeed");
let result = mgr.attach_stream_anchor::<u32>(handle).await;
assert!(
matches!(result, Err(AttachError::AlreadyAttached { .. })),
"concurrent attach must return AlreadyAttached, got {:?}",
result
);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_07_transport_error_propagation() {
let mgr = make_manager();
let mut anchor = mgr.create_anchor::<u32>();
let handle = anchor.handle();
let sender = mgr
.attach_stream_anchor::<u32>(handle)
.await
.expect("attach must succeed");
sender.send(42u32).await.expect("send must succeed");
drop(sender);
let frame1 = tokio::time::timeout(Duration::from_secs(5), anchor.next())
.await
.expect("stream must resolve within 5s");
assert!(
matches!(frame1, Some(Ok(StreamFrame::Item(42u32)))),
"first frame must be Item(42), got {:?}",
frame1
);
let frame2 = tokio::time::timeout(Duration::from_secs(5), anchor.next())
.await
.expect("stream must resolve within 5s");
assert!(
matches!(frame2, Some(Err(StreamError::SenderDropped))),
"second frame must be Err(SenderDropped), got {:?}",
frame2
);
assert!(
anchor.next().await.is_none(),
"stream must be exhausted after SenderDropped"
);
}